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 / scripts / sai-push.service
Author[]Andy Green <andy@warmcat.com> 2021-02-20 13:19 UTC
Committer[]Andy Green <andy@warmcat.com> 2021-02-24 09:57 UTC
Tree1b96543342e0470ca20733365122c6742eab0ec0   Raw Patch
 
resource-manager
resource-manager
diff --git a/CMakeLists.txt b/CMakeLists.txt index 6c4eb43..b856958 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -147,6 +147,7 @@ if (requirements) endif() if (SAI_BUILDER) add_subdirectory(src/builder) + add_subdirectory(src/resource) if (NOT MSVC AND NOT WIN32) add_subdirectory(src/device) add_subdirectory(src/expect) diff --git a/READMEs/README-resource-management.md b/READMEs/README-resource-management.md new file mode 100644 index 0000000..cea813b --- /dev/null +++ b/READMEs/README-resource-management.md @@ -0,0 +1,86 @@ +# Sai Resource Manager + +Once you have many builds running under sai, a serious danger is that +the tests may effectively DDoS any resources they need, like remote +servers. There may be dozens of builders each running builds and the +tests themselves concurrently, at any given moment this may cause +spikes of hundreds or thousands of simultaneous outgoing connections. + +To avoid this, Sai builds in a generic, configurable "resource manager" +in sai-server that is a central place that allocates and queues +requests for leases on amounts of abstract "resources". The actual +resources don't need to be on the server or be physical, so long as +whatever wants to use them participates in the resource management +action via the server. + +![Sai Resource Leasing](../READMEs/sai-resources.png) + +## Configuring well-known resources at sai-server + +The sai-server instance is told about "well-known" resources it is +managing in the sai-server vhost config as a pvo of the +"com-warmcat-sai" protocol named `resources`, eg in +`/etc/sai/server/conf.d/libwebsockets.org` + +``` +... + "ws-protocols": [{ + "com-warmcat-sai": { +... + "resources": "warmcat_conns=32" + } +... +``` + +It takes a comma-separated list in the form of `resource_name=budget`, +the `budget` is the total amount of that resource that may be in use +concurrently. + +## Request architecture + +At the top of the architecture is the centralized resource manager in +sai-server, the sources of the requests for these resources are the +indiviudal tests that require the resources to operate and must block +until they are allocated the resource. + +Since there may be other tests that can run concurrently when these +resources are not available, requesting the resources at the level of +the individual tests can be considerably more efficient. Eg ctest +routinely runs at high levels of test concurrency and so can take +advantage of this. It also allows the amount of the resource being +requested to be tailored to the individual test usage cleanly. + +In the ctest case, tests that require resources define a dependency on +one or more sub-tests that make the resource requests, it is this +dependency that blocks until the resource is made available to it (or +the request fails). + +## `sai-resource` + +`sai-resource` is the interface at the builder. To make sure the +sai-server is scalable, instead of opening its own connection to the +sai-server instance, the sai-builder proxies its existing connection +via a Unix Domain Socket and `sai-resource` connects to that. + +You give it the following arguments to describe what you want + +|argument|example|meaning| +|---|---|---| +|1|warmcat_conns|Well-known resource name| +|2|8|Number of resources needed| +|3|20|Seconds to lease the resources| + +First it creates a random cookie and sends that along with the resource +request to sai-builder, which passes it up to the sai-server to be +queued; the `sai-resource` instance blocks until the resource is +allocated (returning 0) or it times out waiting (returning 1). + +``` +[sai-resource] - [ctest] - [sai-builder] - [sai-server] +``` + +If the resource is allocated, `sai-resource` forks and unblocks the +caller while leaving a fork connected to the Unix Domain Socket, this +allows the `sai-builder` it is connected to to be aware when the test +completed and release the resources at sai-server immediately. + diff --git a/READMEs/README-systemd-nspawn.md b/READMEs/README-systemd-nspawn.md index 1581cf2..45f4e56 100644 --- a/READMEs/README-systemd-nspawn.md +++ b/READMEs/README-systemd-nspawn.md @@ -382,6 +382,8 @@ Container # cd sai && mkdir build && cd build && cmake .. && make && make instal Container # cp ../scripts/sai-builder.service /etc/systemd/system Container # mkdir -p /etc/sai/builder Container # cp ../scripts/builder-conf /etc/sai/builder/conf +Container # mkdir /var/run/com.warmcat.com.saib.logproxy /var/run/com.warmcat.com.saib.resproxy +Container # chown sai /var/run/com.warmcat.com.saib.logproxy /var/run/com.warmcat.com.saib.resproxy Container # vim /etc/sai/builder/conf Container # systemctl enable sai-builder ``` diff --git a/READMEs/sai-resources.png b/READMEs/sai-resources.png new file mode 100644 index 0000000..e741159 Binary files /dev/null and b/READMEs/sai-resources.png differ diff --git a/src/builder/CMakeLists.txt b/src/builder/CMakeLists.txt index 1e234ae..5fb1500 100644 --- a/src/builder/CMakeLists.txt +++ b/src/builder/CMakeLists.txt @@ -10,6 +10,7 @@ set(SRCS b-task.c b-artifacts.c b-logproxy.c + b-refproxy.c ) set(requirements 1) diff --git a/src/builder/b-comms.c b/src/builder/b-comms.c index a0640cb..c38fc52 100644 --- a/src/builder/b-comms.c +++ b/src/builder/b-comms.c @@ -257,6 +257,31 @@ saib_m_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, goto sendify; } + /* + * Any resource requests / relinquishments to process? + */ + + if (spm->resource_req_list.count) { + struct lws_dll2 *d = lws_dll2_get_head(&spm->resource_req_list); + sai_resource_msg_t *resm; + + resm = lws_container_of(d, sai_resource_msg_t, list); + + n = (int)resm->len; + if (*len > resm->len) + *len = resm->len; + memcpy(buf, resm->msg, *len); + + lws_dll2_remove(&resm->list); + free(resm); + + lwsl_notice("%s: forwarding to server %.*s\n", __func__, + (int)(*len), (const char *)buf); + + lws_ss_request_tx(spm->ss); + goto sendify; + } + switch (spm->phase) { case PHASE_BUILDING: case PHASE_IDLE: @@ -279,7 +304,6 @@ saib_m_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, lws_struct_json_serialize_destroy(&js); sp = (sai_plat_t *)builder.sai_plat_owner.head; - lwsl_hexdump_warn(sp, sizeof(*sp)); lwsl_hexdump_notice(start, w); *len = w; diff --git a/src/builder/b-conf.c b/src/builder/b-conf.c index eb9e2c1..f50004e 100644 --- a/src/builder/b-conf.c +++ b/src/builder/b-conf.c @@ -207,20 +207,20 @@ saib_conf_cb(struct lejp_ctx *ctx, char reason) /* * This is the first plat that wants to talk to this server, - * we need to create the logical SS connection + * we need to create the logical SS connection. + * + * The created SS in turn creates a struct sai_plat_server as + * its user object, its CREATING callback in b-comms.c adds + * that to the builder .sai_plat_server_owner */ - if (lws_ss_create(builder.context, 0, &ssi_sai_builder, (void *)ctx, &h, - NULL, NULL)) { + if (lws_ss_create(builder.context, 0, &ssi_sai_builder, + (void *)ctx, &h, NULL, NULL)) { lwsl_err("%s: failed to create secure stream\n", __func__); return -1; } - - - // lws_ss_client_connect(h); - return 0; default: diff --git a/src/builder/b-nspawn.c b/src/builder/b-nspawn.c index fb8514f..1b771f4 100644 --- a/src/builder/b-nspawn.c +++ b/src/builder/b-nspawn.c @@ -203,6 +203,7 @@ ok: static const char * const runscript = "set SAI_INSTANCE_IDX=%d\n" + "set SAI_BUILDER_RESOURCE_PROXY=%s\n" "set SAI_LOGPROXY=%s\n" "set SAI_LOGPROXY_TTY0=%s\n" "set SAI_LOGPROXY_TTY1=%s\n" @@ -217,12 +218,15 @@ static const char * const runscript = "#!/bin/bash -x\n" #if defined(__APPLE__) "export PATH=/opt/homebrew/bin:/usr/local/bin:/usr/bin:/bin:/sbin:/usr/sbin\n" +#else + "export PATH=/usr/local/bin:$PATH\n" #endif "export HOME=%s\n" "export SAI_OVN=%s\n" "export SAI_PROJECT=%s\n" "export SAI_REMOTE_REF=%s\n" "export SAI_INSTANCE_IDX=%d\n" + "export SAI_BUILDER_RESOURCE_PROXY=%s\n" "export SAI_LOGPROXY=%s\n" "export SAI_LOGPROXY_TTY0=%s\n" "export SAI_LOGPROXY_TTY1=%s\n" @@ -240,6 +244,7 @@ saib_spawn(struct sai_nspawn *ns) { struct lws_spawn_piped_info info; char args[290], st[2048], *p; + const char *respath = "unk"; int fd, n; const char * cmd[] = { "/bin/ps", @@ -279,16 +284,24 @@ saib_spawn(struct sai_nspawn *ns) return 1; } + if (builder.sai_plat_server_owner.head) { + struct sai_plat_server *cm = lws_container_of( + builder.sai_plat_server_owner.head, + sai_plat_server_t, list); + + respath = cm->resproxy_path; + } + #if defined(WIN32) n = lws_snprintf(st, sizeof(st), runscript, ns->instance_idx, - ns->slp_control.sockpath, + respath, ns->slp_control.sockpath, ns->slp[0].sockpath, ns->slp[1].sockpath, builder.home, builder.home, ns->fsm.ovname, ns->project_name, ns->task->build); #else n = lws_snprintf(st, sizeof(st), runscript, builder.home, ns->fsm.ovname, ns->project_name, ns->ref, ns->instance_idx, - ns->slp_control.sockpath, + respath, ns->slp_control.sockpath, ns->slp[0].sockpath, ns->slp[1].sockpath, builder.home, ns->task->build); #endif diff --git a/src/builder/b-private.h b/src/builder/b-private.h index 265838c..f80ed28 100644 --- a/src/builder/b-private.h +++ b/src/builder/b-private.h @@ -78,6 +78,11 @@ struct saib_logproxy { int log_channel_idx; }; +struct saib_resproxy { + char sockpath[128]; + struct sai_nspawn *ns; +}; + struct sai_nspawn; struct sai_nspawn { @@ -240,7 +245,14 @@ saib_queue_task_status_update(sai_plat_t *sp, struct sai_plat_server *spm, int rm_rf_cb(const char *dirpath, void *user, struct lws_dir_entry *lde); -extern const struct lws_protocols protocol_logproxy; +extern const struct lws_protocols protocol_logproxy, protocol_resproxy; void * thread_repo(void *d); + +int +saib_create_resproxy_listen_uds(struct lws_context *context, + struct sai_plat_server *spm); + +int +saib_handle_resource_result(struct sai_plat_server *spm, const char *in, size_t len); diff --git a/src/builder/b-refproxy.c b/src/builder/b-refproxy.c new file mode 100644 index 0000000..e7adc15 --- /dev/null +++ b/src/builder/b-refproxy.c @@ -0,0 +1,266 @@ +/* + * sai-builder - resproxy + * + * Copyright (C) 2019 - 2021 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 + * + * We create a listening UDS socket for each server that we accept jobs from. + * sai-resource can later connect to it (the path is passed in by env var) + * and do blocking requests for resource leases, proxied by us to the server + * we're already connected to. + * + * This code handles the accepted UDS connections that connected to one of these + * resource proxy UDS listeners, and performs the proxy action to the associated + * server using the existing link the builder already has to it. + * + * The client sends us a well-formed JSON we just forward to the server. We + * know there's a cookie in there we stash a copy of in the pss, otherwise we + * don't understand the JSON. + * + * The server responds with a JSON (or the client times out waiting) that again + * we don't understand except the cookie, to find out which pss it belongs to. + * + * If we can't match the cookie, or the client hangs up, we send a JSON with + * just the schema and the cookie to explicitly relinquish the remaining lease, + * that stops us always waiting out the whole lease period when we just used + * a small part of it. + */ + +#include <libwebsockets.h> +#include <string.h> +#include <signal.h> +#include "b-private.h" + +struct rppss { + struct lws *wsi; + lws_dll2_t list; /* builder's list of resource pss' */ + + char cookie[65]; + + char *response; + unsigned int response_len; + + uint8_t accepted:1; +}; + +static struct rppss * +resproxy_find_by_cookie(struct sai_plat_server *spm, const char *c, size_t clen) +{ + lws_start_foreach_dll(lws_dll2_t *, p, spm->resource_pss_list.head) { + struct rppss *pss = lws_container_of(p, struct rppss, list); + + if (!strncmp(pss->cookie, c, clen)) + return pss; + + } lws_end_foreach_dll(p); + + return NULL; +} + +static int +saib_queue_yield_message(struct sai_plat_server *spm, const char *c, size_t len) +{ + sai_resource_msg_t *m = + malloc(sizeof(sai_resource_msg_t) + LWS_PRE + 196); + + if (!m) + return -1; + + memset(m, 0, sizeof(*m)); + m->msg = ((char *)&m[1]) + LWS_PRE; + /* + * We just send the cookie to relinquish the leased resources + */ + m->len = (size_t)lws_snprintf((char *)m->msg, 196, + "{\"schema\":\"com-warmcat-sai-resource\"," + "\"cookie\":\"%.*s\"}", (int)len, c); + + lwsl_notice("%s: %s\n", __func__, m->msg); + + lws_dll2_add_tail(&m->list, &spm->resource_req_list); + lws_ss_request_tx(spm->ss); + + return 0; +} + +int +saib_handle_resource_result(struct sai_plat_server *spm, const char *in, size_t len) +{ + struct rppss *pss; + const char *p; + size_t al; + + p = lws_json_simple_find((const char *)in, len, "\"cookie\":", &al); + if (!p) { + /* seems malformed */ + lwsl_warn("%s: server sent JSON without a cookie\n", __func__); + return -1; + } + + pss = resproxy_find_by_cookie(spm, p, al); + if (!pss) { + lwsl_warn("%s: the requestor left before the response\n", + __func__); + /* + * Explicit yield, in case the acceptance raced the client + * closing... if it was telling us we can't have it, the server + * will ignore the yield since no lease with this cookie on + * record there + */ + saib_queue_yield_message(spm, p, al); + + return 1; /* the requestor has gone away */ + } + + /* + * We want to proxy back the server's response to the client that + * we identified owns the cookie, we don't need to understand the + * response any further ourselves + */ + + pss->response_len = (unsigned int)len; + pss->response = malloc(len + LWS_PRE); + if (!pss->response) + return 0; + + memcpy(pss->response + LWS_PRE, in, len); + + lwsl_notice("%s: queuing server response %.*s\n", __func__, (int)len, in); + + /* + * Mark us as having been handled, so we will have to yield it + * when we close. If this gets lost somewhere the lease period expiring + * will autorecover it at the server side, so no worries... + */ + pss->accepted = 1; + lws_callback_on_writable(pss->wsi); + + return 0; /* we will try to pass it on */ +} + +static int +callback_resproxy(struct lws *wsi, enum lws_callback_reasons reason, + void *user, void *in, size_t len) +{ + struct sai_plat_server *spm = lws_vhost_user(lws_get_vhost(wsi)); + struct rppss *pss = (struct rppss *)user; + sai_resource_msg_t *m; + const char *p; + size_t al; + + switch (reason) { + + case LWS_CALLBACK_RAW_CLOSE: + + lwsl_notice("%s: CLOSE\n", __func__); + /* + * We can close at any time regardless of what the server's + * doing... delist and destroy our request + */ + + /* stop tracking the pss for this, we're closing */ + lws_dll2_remove(&pss->list); + + if (pss->response) { + free(pss->response); + pss->response = NULL; + } + + if (pss->accepted) + saib_queue_yield_message(spm, pss->cookie, + strlen(pss->cookie)); + break; + + case LWS_CALLBACK_RAW_RX: + + if (pss->wsi) + /* only request once */ + return 0; + + pss->wsi = wsi; + + /* + * We get sent something like this + * + * { + * "schema":"com-warmcat-sai-resource", + * "resname":"warmcat_conns", + * "cookie":"xxx", + * "amount":8, + * "lease":20 + * } + * + * We keep a copy of the cookie in the pss, and add the pss to + * an owner in the server until it is closed. + */ + + p = lws_json_simple_find((const char *)in, len, + "\"cookie\":", &al); + if (!p) + return -1; + + lws_strnncpy(pss->cookie, p, al, sizeof(pss->cookie)); + + m = malloc(sizeof(*m) + LWS_PRE + len); + memset(m, 0, sizeof(*m)); + m->msg = ((const char *)&m[1] + LWS_PRE); + memcpy((char *)m->msg, in, len); + m->len = len; + + lws_dll2_add_tail(&m->list, &spm->resource_req_list); + + lws_dll2_add_tail(&pss->list, &spm->resource_pss_list); + + /* try to schedule a write */ + lws_ss_request_tx(spm->ss); + break; + + case LWS_CALLBACK_RAW_WRITEABLE: + if (pss->response) { + lwsl_notice("%s: WRITABLE issue response %.*s\n", + __func__, (int)pss->response_len, + pss->response + LWS_PRE); + if (lws_write(wsi, (uint8_t *)pss->response + LWS_PRE, + (unsigned int)pss->response_len, + LWS_WRITE_RAW) != (int)pss->response_len) + return -1; + + free(pss->response); + pss->response = NULL; + + /* + * We stay up until the client goes away, triggering + * explicitly giving up the lease + */ + + break; + } + break; + + default: + break; + } + + return 0; +} + +const struct lws_protocols protocol_resproxy = { + "protocol-resproxy", + callback_resproxy, + sizeof(struct rppss), + 2048, 2048, NULL, 0 +}; diff --git a/src/builder/b-sai.c b/src/builder/b-sai.c index 92566b5..9b8a66b 100644 --- a/src/builder/b-sai.c +++ b/src/builder/b-sai.c @@ -114,6 +114,7 @@ static const char * const default_ss_policy = static const struct lws_protocols *pprotocols[] = { &protocol_stdxxx, &protocol_logproxy, + &protocol_resproxy, &lws_openmetrics_export_protocols[LWSOMPROIDX_PROX_WS_CLIENT], NULL }; @@ -140,18 +141,26 @@ pvo1a = { "ws-server-uri", /* protocol name we belong to on this vhost */ "ok" /* set at runtime from conf */ }, -pvo1 = { +pvo1 = { /* starting point for metrics proxy */ NULL, /* "next" pvo linked-list */ &pvo1a, /* "child" pvo linked-list */ "lws-openmetrics-prox-client", /* protocol name we belong to on this vhost */ "ok" /* ignored */ }, -pvo = { + +pvo = { /* starting point for logproxy */ NULL, /* "next" pvo linked-list */ NULL, /* "child" pvo linked-list */ "protocol-logproxy", /* protocol name we belong to on this vhost */ "ok" /* ignored */ -}; +}, + +pvo_resproxy = { /* starting point for resproxy */ + NULL, /* "next" pvo linked-list */ + NULL, /* "child" pvo linked-list */ + "protocol-resproxy", /* protocol name we belong to on this vhost */ + "ok" /* ignored */ + };; static int saib_create_listen_uds(struct lws_context *context, struct saib_logproxy *lp) @@ -175,7 +184,44 @@ saib_create_listen_uds(struct lws_context *context, struct saib_logproxy *lp) unlink(lp->sockpath); #endif - lwsl_notice("%s: %s\n", __func__, lp->sockpath); + lwsl_notice("%s: %s.%s\n", __func__, info.vhost_name, lp->sockpath); + + if (!lws_create_vhost(context, &info)) { + lwsl_notice("%s: failed to create vh %s\n", __func__, + info.vhost_name); + return -1; + } + + return 0; +} + +/* + * We create one of these per server we connected to + */ + +int +saib_create_resproxy_listen_uds(struct lws_context *context, + struct sai_plat_server *spm) +{ + struct lws_context_creation_info info; + + memset(&info, 0, sizeof(info)); + + info.vhost_name = pv; + pv += lws_snprintf(pv, sizeof(vhnames) - (size_t)(pv - vhnames), + "resproxy.%d", spm->index) + 1; + info.options = LWS_SERVER_OPTION_ADOPT_APPLY_LISTEN_ACCEPT_CONFIG | + LWS_SERVER_OPTION_UNIX_SOCK; + + info.iface = spm->resproxy_path; + info.listen_accept_role = "raw-skt"; + info.listen_accept_protocol = "protocol-resproxy"; + info.user = spm; + info.pvo = &pvo_resproxy; + info.pprotocols = pprotocols; + + lwsl_notice("%s: Created resproxy %s.%s\n", __func__, info.vhost_name, + spm->resproxy_path); if (!lws_create_vhost(context, &info)) { lwsl_notice("%s: failed to create vh %s\n", __func__, @@ -298,6 +344,31 @@ app_system_state_nf(lws_state_manager_t *mgr, lws_state_notify_link_t *link, } lws_end_foreach_dll_safe(np, np1); } lws_end_foreach_dll_safe(mp, mp1); + /* + * Create the resource proxy listeners, one per server link + */ + + lwsl_notice("%s: creating resource proxy listeners\n", __func__); + + lws_start_foreach_dll(struct lws_dll2 *, pxx, + builder.sai_plat_server_owner.head) { + struct sai_plat_server *spm = lws_container_of(pxx, sai_plat_server_t, list); + + lws_snprintf(spm->resproxy_path, sizeof(spm->resproxy_path), + #if defined(__linux__) + UDS_PATHNAME_RESPROXY".%d", + #else + UDS_PATHNAME_RESPROXY"/%d", + #endif + spm->index); + + lwsl_notice("%s: creating %s\n", __func__, spm->resproxy_path); + + saib_create_resproxy_listen_uds(builder.context, spm); + + } lws_end_foreach_dll(pxx); + + break; } @@ -394,7 +465,8 @@ int main(int argc, const char **argv) return 1; } - lwsl_notice("%s: parsed %s %s %s\n", __func__, builder.metrics_path, builder.metrics_uri, builder.metrics_secret); +// lwsl_notice("%s: parsed %s %s %s\n", __func__, builder.metrics_path, +// builder.metrics_uri, builder.metrics_secret); /* * We need to sample the true uid / gid we should use inside @@ -476,6 +548,7 @@ int main(int argc, const char **argv) pthread_mutex_init(&builder.mi.mut, NULL); pthread_cond_init(&builder.mi.cond, NULL); + /* * Our approach is to split off a thread to do the git remote handling * in a serialized way blocking the related threadpool threads until it diff --git a/src/builder/b-task.c b/src/builder/b-task.c index 817a9f6..d753427 100644 --- a/src/builder/b-task.c +++ b/src/builder/b-task.c @@ -36,11 +36,13 @@ static char csep = '\\'; const lws_struct_map_t lsm_schema_map_m_to_b[] = { LSM_SCHEMA (sai_task_t, NULL, lsm_task, "com-warmcat-sai-ta"), LSM_SCHEMA (sai_cancel_t, NULL, lsm_task_cancel, "com.warmcat.sai.taskcan"), + LSM_SCHEMA (sai_resource_t, NULL, lsm_resource, "com-warmcat-sai-resource") }; enum { SAIB_RX_TASK_ALLOCATION, SAIB_RX_TASK_CANCEL, + SAIB_RX_RESOURCE_REPLY, }; static const char * const nsstates[] = { @@ -411,6 +413,7 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len) struct lws_threadpool_task_args tpa; sai_plat_t *sp = NULL; struct sai_nspawn *ns; + sai_resource_t *reso; struct lejp_ctx ctx; lws_struct_args_t a; sai_cancel_t *can; @@ -433,8 +436,8 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len) lws_struct_json_init_parse(&ctx, NULL, &a); m = lejp_parse(&ctx, (uint8_t *)in, (int)len); if (m < 0) { - printf("%.*s\n", (int)len, (const char *)in); - lwsl_notice("%s: builder rx JSON decode failed '%s'\n", + lwsl_hexdump_err(in, len); + lwsl_err("%s: builder rx JSON decode failed '%s'\n", __func__, lejp_error_to_string(m)); return m; } @@ -739,6 +742,15 @@ saib_ws_json_rx_builder(struct sai_plat_server *spm, const void *in, size_t len) } lws_end_foreach_dll_safe(mp, mp1); break; + case SAIB_RX_RESOURCE_REPLY: + reso = (sai_resource_t *)a.dest; + + lwsl_notice("%s: RESOURCE_REPLY: cookie %s\n", + __func__, reso->cookie); + + saib_handle_resource_result(spm, in, len); + break; + default: break; } diff --git a/src/common/include/private.h b/src/common/include/private.h index 22bfbc6..3d114c8 100644 --- a/src/common/include/private.h +++ b/src/common/include/private.h @@ -1,7 +1,7 @@ /* * Sai - ./src/common/include/private.h * - * Copyright (C) 2019-2020 Andy Green <andy@warmcat.com> + * Copyright (C) 2019 - 2021 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 @@ -25,8 +25,10 @@ #if defined(__linux__) #define UDS_PATHNAME_LOGPROXY "@com.warmcat.com.saib.logproxy" +#define UDS_PATHNAME_RESPROXY "@com.warmcat.com.saib.resproxy" #else #define UDS_PATHNAME_LOGPROXY "/var/run/com.warmcat.com.saib.logproxy" +#define UDS_PATHNAME_RESPROXY "/var/run/com.warmcat.com.saib.resproxy" #endif struct sai_plat; @@ -160,6 +162,58 @@ typedef struct { char sent_json; } sai_artifact_t; +/* communication part of resource allocation requests */ + +typedef struct { + const char *resname; + const char *cookie; + unsigned int amount; + unsigned int lease; +} sai_resource_t; + +typedef struct { + lws_dll2_t list_resource_wellknown; + lws_dll2_t list_resource_queued_leased; + lws_dll2_t list_pss; + + lws_sorted_usec_list_t sul_expiry; + + const char *cookie; + + time_t requested_since_time; + time_t allocated_since_time; + unsigned int amount; + unsigned int lease_secs; + + /* cookie is overallocated */ +} sai_resource_requisition_t; + +typedef struct { + lws_dll2_t list; + + struct lws_context *cx; + + /* any related resources listed here so we can get this object */ + lws_dll2_owner_t owner; /* sai_resource_requisition_t */ + /* queue for pending requests on this resource */ + lws_dll2_owner_t owner_queued; /* sai_resource_requisition_t */ + /* list of allocated leases */ + lws_dll2_owner_t owner_leased; /* sai_resource_requisition_t */ + + const char *name; + long budget; + long allocated; + + /* name is overallocated */ +} sai_resource_wellknown_t; + +typedef struct { + lws_dll2_t list; + const char *msg; + size_t len; + /* msg is overallocated */ +} sai_resource_msg_t; + struct sai_plat; /* @@ -173,6 +227,10 @@ typedef struct sai_plat_server { lws_dll2_t list; lws_dll2_owner_t rejection_list; + lws_dll2_owner_t resource_req_list; /* sai_resource_msg_t */ + lws_dll2_owner_t resource_pss_list; /* so we can find the cookie */ + + char resproxy_path[128]; const char *url; const char *name; @@ -282,7 +340,8 @@ extern const lws_struct_map_t lsm_task_cancel[1], lsm_schema_json_map_can[1], lsm_schema_json_map_task[1], - lsm_schema_json_map_event[1] + lsm_schema_json_map_event[1], + lsm_resource[4] ; extern const lws_ss_info_t ssi_said_logproxy; diff --git a/src/common/ss-client-logproxy.c b/src/common/ss-client-logproxy.c index 40c550c..b880163 100644 --- a/src/common/ss-client-logproxy.c +++ b/src/common/ss-client-logproxy.c @@ -41,6 +41,7 @@ typedef struct said_logproxy { lws_dll2_owner_t logs; struct lws_ss_handle *ss; void *opaque_data; + char pending; } said_logproxy_t; static saicom_drain_cb cb; diff --git a/src/common/struct-metadata.c b/src/common/struct-metadata.c index 67ff0ed..b96f3c4 100644 --- a/src/common/struct-metadata.c +++ b/src/common/struct-metadata.c @@ -1,7 +1,7 @@ /* - * Sai server - ./src/server/struct-metadata.c + * Sai server - ./src/common/struct-metadata.c * - * Copyright (C) 2019 - 2020 Andy Green <andy@warmcat.com> + * Copyright (C) 2019 - 2021 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 @@ -156,6 +156,18 @@ const lws_struct_map_t lsm_schema_sq3_map_log[] = { LSM_SCHEMA_DLL2 (sai_log_t, list, NULL, lsm_log, "logs"), }; +const lws_struct_map_t lsm_resource[] = { + LSM_STRING_PTR (sai_resource_t, resname, "resname"), + LSM_STRING_PTR (sai_resource_t, cookie, "cookie"), + LSM_UNSIGNED (sai_resource_t, amount, "amount"), + LSM_UNSIGNED (sai_resource_t, lease, "lease"), +}; + +const lws_struct_map_t lsm_schema_json_map_resource[] = { + LSM_SCHEMA (sai_resource_t, NULL, lsm_resource, + "com-warmcat-sai-resource"), +}; + /* * Artifacts live in their own table in the event-specific db, and refer back to * a task uuid. The build decides how many artifacts exist for a task (the diff --git a/src/resource/CMakeLists.txt b/src/resource/CMakeLists.txt new file mode 100644 index 0000000..5eb6cb6 --- /dev/null +++ b/src/resource/CMakeLists.txt @@ -0,0 +1,41 @@ +set(SUB "sai-resource") +set(CPACK_DEBIAN_BUILDER_PACKAGE_NAME "sai-resource") + +set(SRCS + r-sai.c + r-comms.c +) + +set(requirements 1) +require_lws_config(LWS_WITH_CLIENT 1 requirements) +require_lws_config(LWS_WITH_SPAWN 1 requirements) +require_lws_config(LWS_WITH_STRUCT_JSON 1 requirements) +require_lws_config(LWS_WITH_SECURE_STREAMS 1 requirements) +require_lws_config(LWS_WITH_DIR 1 requirements) + +if (requirements AND NOT MSVC) + + add_executable(${SUB} ${SRCS}) + if (APPLE) + set_property(TARGET sai-resource PROPERTY MACOSX_RPATH YES) + endif() + if (SAI_LWS_INC_PATH) + target_include_directories(${SUB} PRIVATE ${SAI_LWS_INC_PATH}) + endif() + + target_link_libraries(${SUB} websockets ${SAI_LWS_LIB_PATH}) + + message("LWS_OPENSSL_LIBRARIES ${SUB} '${LWS_OPENSSL_LIBRARIES}'") + if (LWS_OPENSSL_LIBRARIES) + target_link_libraries(${SUB} ${LWS_OPENSSL_LIBRARIES}) + endif() + + if (MSVC OR WIN32) + target_link_libraries(${SUB} ws2_32.lib userenv.lib psapi.lib iphlpapi.lib) + endif() + + install(TARGETS ${SUB} + RUNTIME DESTINATION "${BIN_DIR}" COMPONENT resource) + +endif(requirements AND NOT MSVC) +include(CPack) diff --git a/src/resource/r-comms.c b/src/resource/r-comms.c new file mode 100644 index 0000000..8e5dbca --- /dev/null +++ b/src/resource/r-comms.c @@ -0,0 +1,137 @@ +/* + * sai-resource + * + * Copyright (C) 2019 - 2021 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 + */ + +#include <libwebsockets.h> +#include <string.h> +#include <signal.h> + +#include "../common/include/private.h" +#include "r-private.h" + +typedef struct sair_resource { + struct lws_ss_handle *ss; + void *opaque_data; +} sair_resource_t; + +static lws_ss_state_return_t +sair_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags) +{ + /* + * Possibly: 1) nothing comes and we get timed out by whatever spawned + * us (usually ctest), 2) we get the lease acknowledgement from + * sai-server and can exit successfully, or 3) we get an explicit deny + * from sai-server (eg, unknown well-known resource name). + * + * He sends us a JSON with an optional "error" description member, + * "result" being 0 means it went OK, anything else we failed. + */ + + lwsl_hexdump_notice(buf, len); + + return LWSSSSRET_OK; +} + +static lws_ss_state_return_t +sair_lp_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, size_t *len, + int *flags) +{ +// sair_resource_t *lp = (sair_resource_t *)userobj; + + if (!msg[0]) + return LWSSSSRET_TX_DONT_SEND; /* nothing to send */ + + *flags = LWSSS_FLAG_SOM | LWSSS_FLAG_EOM; + if (*len > strlen(msg)) + *len = strlen(msg); + + memcpy(buf, msg, *len); + msg[0] = '\0'; + + lwsl_notice("%s: issuing requisition %s\n", __func__, msg); + + return LWSSSSRET_OK; +} + +static lws_ss_state_return_t +sair_lp_state(void *userobj, void *sh, lws_ss_constate_t state, + lws_ss_tx_ordinal_t ack) +{ + sair_resource_t *lp = (sair_resource_t *)userobj; + + lwsl_info("%s: %s, ord 0x%x\n", __func__, lws_ss_state_name((int)state), + (unsigned int)ack); + + switch (state) { + + case LWSSSCS_CONNECTED: + lws_ss_request_tx(lp->ss); + break; + + case LWSSSCS_ALL_RETRIES_FAILED: + case LWSSSCS_DISCONNECTED: + /* + * That's it for us + */ + interrupted = 1; + break; + + default: + break; + } + + return LWSSSSRET_OK; +} + +static const lws_ss_info_t ssi_sair_resource = { + .handle_offset = offsetof(sair_resource_t, ss), + .opaque_user_data_offset = offsetof(sair_resource_t, opaque_data), + .rx = sair_lp_rx, + .tx = sair_lp_tx, + .state = sair_lp_state, + .user_alloc = sizeof(sair_resource_t), + .streamtype = "resproxy" +}; + +struct lws_ss_handle * +sair_ss_from_env(struct lws_context *context, const char *env_name) +{ + char *e = getenv(env_name); + struct lws_ss_handle *h; + + if (!e) { + lwsl_notice("%s: no env var %s\n", __func__, env_name); + + return NULL; + } + + lwsl_notice("%s: connecting to %s\n", __func__, e); + + if (lws_ss_create(context, 0, &ssi_sair_resource, 0, &h, NULL, NULL)) { + lwsl_err("%s: failed to create SS for %s\n", __func__, env_name); + + return NULL; + } + + lws_ss_set_metadata(h, "sockpath", e, strlen(e)); + lws_ss_client_connect(h); + + return h; +} diff --git a/src/resource/r-private.h b/src/resource/r-private.h new file mode 100644 index 0000000..3575b2c --- /dev/null +++ b/src/resource/r-private.h @@ -0,0 +1,26 @@ +/* + * sai-resource + * + * Copyright (C) 2019 - 2021 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 + */ + +extern char msg[]; +extern int interrupted; + +struct lws_ss_handle * +sair_ss_from_env(struct lws_context *context, const char *env_name); diff --git a/src/resource/r-sai.c b/src/resource/r-sai.c new file mode 100644 index 0000000..dcbd33d --- /dev/null +++ b/src/resource/r-sai.c @@ -0,0 +1,167 @@ +/* + * sai-resource + * + * Copyright (C) 2019 - 2021 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 + * + * sai-resource : helper to make blocking requests for resources from + * remote sai-server, via an existing sai-builder connection + */ + +#include <libwebsockets.h> +#include <string.h> +#include <signal.h> + +#include "r-private.h" + +struct lws_ss_handle *ssh; +static lws_state_notify_link_t nl; +static struct lws_vhost *vh; +int interrupted, bad = 1; +char msg[1024], cookie[165]; + +static const char * const default_ss_policy = + "{" + "\"s\": [" + /* + * Unix Domain Socket connections to sai-builder resource proxy + */ + "{\"resproxy\": {" + "\"endpoint\":" "\"+${sockpath}\"," + "\"protocol\":" "\"raw\"," + "\"opportunistic\":" "true," + "\"metadata\": [" + "{\"sockpath\": \"\"}" + "]" + "}}" + "]}" +; + +static lws_state_notify_link_t * const app_notifier_list[] = { + &nl, NULL +}; + +static int +app_system_state_nf(lws_state_manager_t *mgr, lws_state_notify_link_t *link, + int current, int target) +{ + struct lws_context *context = lws_system_context_from_system_mgr(mgr); + + switch (target) { + + case LWS_SYSTATE_OPERATIONAL: + if (current != LWS_SYSTATE_OPERATIONAL) + break; + + ssh = sair_ss_from_env(context, "SAI_BUILDER_RESOURCE_PROXY"); + if (!ssh) + goto bail; + break; + } + + return 0; + +bail: + interrupted = 1; + + return 1; +} + +void sigint_handler(int sig) +{ + interrupted = 1; +} + +int main(int argc, const char **argv) +{ + int logs = LLL_USER | LLL_ERR | LLL_WARN | LLL_NOTICE; + struct lws_context_creation_info info; + struct lws_context *context; + const char *p; + + if (argc < 5) { + lwsl_err("%s: Usage: sai-resource <name> <amount> <lease-secs> <cookie>\n", __func__); + return 1; + } + + memset(&info, 0, sizeof info); + + if ((p = lws_cmdline_option(argc, argv, "-d"))) + logs = atoi(p); + + lws_set_log_level(logs, NULL); + + lwsl_user("Sai Resource - " + "Copyright (C) 2019-2021 Andy Green <andy@warmcat.com>\n"); + + info.port = CONTEXT_PORT_NO_LISTEN; + + info.pss_policies_json = default_ss_policy; + info.pt_serv_buf_size = 1 * 1024; + info.options = LWS_SERVER_OPTION_EXPLICIT_VHOSTS; + + signal(SIGINT, sigint_handler); + info.fd_limit_per_thread = 1 + 8 + 1; + + /* hook up our lws_system state notifier */ + + nl.name = "sai-resource"; + nl.notify_cb = app_system_state_nf; + info.register_notifier_list = app_notifier_list; + + /* create the lws context */ + + context = lws_create_context(&info); + if (!context) { + lwsl_err("lws init failed\n"); + return 1; + } + + lws_strncpy(cookie, argv[4], sizeof(cookie)); + + /* + * args: 1) resource well-known name + * 2) amount of resource requested + * 3) requested lease time (in seconds) + * 4) uid string for this lease + */ + + lws_snprintf(msg, sizeof(msg), + "{\"schema\":\"com-warmcat-sai-resource\"," + "\"resname\":\"%s\"," + "\"cookie\":\"%s\"," + "\"amount\":%u," + "\"lease\":%u}", argv[1], cookie, + (unsigned int)atol(argv[2]), + (unsigned int)atol(argv[3])); + + /* ... and our vhost... */ + + vh = lws_create_vhost(context, &info); + if (!vh) { + lwsl_err("%s: Failed to create vhost\n", __func__); + goto bail1; + } + + while (!lws_service(context, 0) && !interrupted) + ; + +bail1: + lws_context_destroy(context); + + return bad; +} diff --git a/src/server/CMakeLists.txt b/src/server/CMakeLists.txt index b256075..74e4e8d 100644 --- a/src/server/CMakeLists.txt +++ b/src/server/CMakeLists.txt @@ -11,6 +11,7 @@ set(SRCS s-task.c s-central.c s-websrv.c + s-resource.c ) set(requirements 1) diff --git a/src/server/s-comms.c b/src/server/s-comms.c index f132d6d..766b887 100644 --- a/src/server/s-comms.c +++ b/src/server/s-comms.c @@ -267,7 +267,7 @@ sais_all_browser_on_writable(struct vhd *vhd) #endif static int -sai_destroy_builder(struct lws_dll2 *d, void *user) +sai_detach_builder(struct lws_dll2 *d, void *user) { // saib_t *b = lws_container_of(d, saib_t, c.builder_list); @@ -276,16 +276,50 @@ sai_destroy_builder(struct lws_dll2 *d, void *user) return 0; } +static int +sai_detach_resource(struct lws_dll2 *d, void *user) +{ + lws_dll2_remove(d); + + return 0; +} + +static int +sai_destroy_resource_wellknown(struct lws_dll2 *d, void *user) +{ + sai_resource_wellknown_t *rwk = + lws_container_of(d, sai_resource_wellknown_t, list); + + /* + * Just detach everything listed on this well-known resource... + * everything listed here is ultimately owned by a pss and will be + * destroyed when that goes down + */ + + lws_dll2_foreach_safe(&rwk->owner_queued, NULL, sai_detach_resource); + lws_dll2_foreach_safe(&rwk->owner_leased, NULL, sai_detach_resource); + + lws_dll2_remove(d); + + free(rwk); + + return 0; +} + static void sais_server_destroy(struct vhd *vhd, sais_t *server) { lwsl_notice("%s: server %p\n", __func__, server); if (server) - lws_dll2_foreach_safe(&server->builder_owner, NULL, sai_destroy_builder); + lws_dll2_foreach_safe(&server->builder_owner, NULL, + sai_detach_builder); sais_event_db_close_all_now(vhd); lws_struct_sq3_close(&server->pdb); + + lws_dll2_foreach_safe(&server->resource_wellknown_owner, NULL, + sai_destroy_resource_wellknown); } typedef enum { @@ -338,6 +372,7 @@ callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, *end = &buf[sizeof(buf) - LWS_PRE - 1]; struct pss *pss = (struct pss *)user; sai_http_murl_t mu = SHMUT_NONE; + const char *pvo_resources; int n; (void)end; @@ -365,6 +400,55 @@ callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, return -1; } + /* + * Create the listed well-known resources to be managed by the + * sai-server for the builders + */ + + if (!lws_pvo_get_str(in, "resources", &pvo_resources)) { + sai_resource_wellknown_t *wk; + struct lws_tokenize ts; + char wkname[32]; + + wkname[0] = '\0'; + lws_tokenize_init(&ts, pvo_resources, + LWS_TOKENIZE_F_MINUS_NONTERM); + do { + + ts.e = (int8_t)lws_tokenize(&ts); + switch (ts.e) { + case LWS_TOKZE_TOKEN_NAME_EQUALS: + lws_strnncpy(wkname, ts.token, ts.token_len, + sizeof(wkname)); + break; + case LWS_TOKZE_INTEGER: + + /* + * Create a new well-known resource + */ + + wk = malloc(sizeof(*wk) + strlen(wkname) + 1); + if (!wk) + return -1; + + memset(wk, 0, sizeof(*wk)); + wk->cx = lws_get_context(wsi); + wk->name = (const char *)&wk[1]; + memcpy((char *)wk->name, wkname, + strlen(wkname) + 1); + wk->budget = atol(ts.token); + lwsl_notice("%s: well-known resource '%s' " + "initialized to %ld\n", __func__, + wk->name, wk->budget); + lws_dll2_add_tail(&wk->list, &vhd->server. + resource_wellknown_owner); + break; + default: + break; + } + } while (ts.e > 0); + } + lws_snprintf((char *)buf, sizeof(buf), "%s-events.sqlite3", vhd->sqlite3_path_lhs); @@ -635,6 +719,8 @@ callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, } lws_end_foreach_dll_safe(p, p1); + sais_resource_wellknown_remove_pss(&pss->vhd->server, pss); + if (pss->blob_artifact) { sqlite3_blob_close(pss->blob_artifact); pss->blob_artifact = NULL; diff --git a/src/server/s-private.h b/src/server/s-private.h index 38f707d..68101ec 100644 --- a/src/server/s-private.h +++ b/src/server/s-private.h @@ -29,8 +29,11 @@ struct sai_plat; typedef struct sai_platm { - struct lws_dll2_owner builder_owner; - struct lws_dll2_owner subs_owner; + lws_dll2_owner_t builder_owner; + lws_dll2_owner_t subs_owner; + + /* the list of well-known, configured resources */ + lws_dll2_owner_t resource_wellknown_owner; /* sai_resource_wellknown_t */ sqlite3 *pdb; sqlite3 *pdb_auth; @@ -97,6 +100,12 @@ struct pss { 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 aft_owner; /* for statefully spooling artifact info */ + lws_dll2_owner_t res_owner; /* sai_resource_requisition_t + * owner of resource objects related + * to this pss */ + lws_dll2_owner_t res_pending_reply_owner; /* sai_resource_msg_t + * resource JSON return + * messages to builder */ lws_struct_args_t a; union { @@ -261,3 +270,21 @@ sais_websrv_broadcast(struct lws_ss_handle *hsrv, const char *str, size_t len); int sql3_get_integer_cb(void *user, int cols, char **values, char **name); + +sai_resource_wellknown_t * +sais_resource_wellknown_by_name(sais_t *sais, const char *name); + +void +sais_resource_wellknown_remove_pss(sais_t *sais, struct pss *pss); + +void +sais_resource_check_if_can_accept_queued(sai_resource_wellknown_t *wk); + +sai_resource_requisition_t * +sais_resource_lookup_lease_by_cookie(sais_t *sais, const char *cookie); + +void +sais_resource_destroy_queued_by_cookie(sais_t *sais, const char *cookie); + +void +sais_resource_rr_destroy(sai_resource_requisition_t *rr); diff --git a/src/server/s-resource.c b/src/server/s-resource.c new file mode 100644 index 0000000..0ebb46c --- /dev/null +++ b/src/server/s-resource.c @@ -0,0 +1,240 @@ +/* + * Sai server - ./src/server/s-resource.c + * + * Copyright (C) 2019 - 2021 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 + * + * Serverside resource allocation implementation + */ + +#include <libwebsockets.h> +#include <string.h> +#include <signal.h> +#include <time.h> +#include <stdio.h> +#include <fcntl.h> +#include <assert.h> + +#include "s-private.h" + +sai_resource_wellknown_t * +sais_resource_wellknown_by_name(sais_t *sais, const char *name) +{ + lws_start_foreach_dll(struct lws_dll2 *, p, + sais->resource_wellknown_owner.head) { + sai_resource_wellknown_t *wk; + + wk = lws_container_of(p, sai_resource_wellknown_t, list); + + if (!strcmp(name, wk->name)) + return wk; + + } lws_end_foreach_dll(p); + + return NULL; +} + +sai_resource_requisition_t * +sais_resource_lookup_lease_by_cookie(sais_t *sais, const char *cookie) +{ + lws_start_foreach_dll(struct lws_dll2 *, p, + sais->resource_wellknown_owner.head) { + sai_resource_wellknown_t *wk; + + wk = lws_container_of(p, sai_resource_wellknown_t, list); + + lws_start_foreach_dll(struct lws_dll2 *, p1, + wk->owner_leased.head) { + sai_resource_requisition_t *rr = lws_container_of(p1, + sai_resource_requisition_t, + list_resource_queued_leased); + + if (!strcmp(cookie, rr->cookie)) + return rr; + + } lws_end_foreach_dll(p1); + + } lws_end_foreach_dll(p); + + return NULL; +} + +void +sais_resource_rr_destroy(sai_resource_requisition_t *rr) +{ + sai_resource_wellknown_t *wk; + + wk = lws_container_of(rr->list_resource_wellknown.owner, + sai_resource_wellknown_t, owner); + + if (rr->allocated_since_time) { + assert((time_t)wk->allocated >= (time_t)rr->amount); + wk->allocated = wk->allocated - (long)rr->amount; + lwsl_notice("%s: freeing lease -> %lu/%lu\n", __func__, + wk->allocated, wk->budget); + } + + lws_sul_cancel(&rr->sul_expiry); + + lws_dll2_remove(&rr->list_resource_wellknown); + lws_dll2_remove(&rr->list_resource_queued_leased); + lws_dll2_remove(&rr->list_pss); + + free(rr); + + sais_resource_check_if_can_accept_queued(wk); +} + +/* + * Eradicate any request, from its cookie + */ + +void +sais_resource_destroy_queued_by_cookie(sais_t *sais, const char *cookie) +{ + lws_start_foreach_dll(struct lws_dll2 *, p, + sais->resource_wellknown_owner.head) { + sai_resource_wellknown_t *wk; + + wk = lws_container_of(p, sai_resource_wellknown_t, list); + + lws_start_foreach_dll(struct lws_dll2 *, p1, + wk->owner_queued.head) { + sai_resource_requisition_t *rr = lws_container_of(p1, + sai_resource_requisition_t, + list_resource_queued_leased); + + if (!strcmp(cookie, rr->cookie)) { + sais_resource_rr_destroy(rr); + return; + } + + } lws_end_foreach_dll(p1); + + } lws_end_foreach_dll(p); +} + +static void +sais_res_expiry(lws_sorted_usec_list_t *sul) +{ + sai_resource_requisition_t *rr = lws_container_of(sul, + sai_resource_requisition_t, sul_expiry); + + lwsl_notice("%s: lease expired\n", __func__); + + sais_resource_rr_destroy(rr); +} + +/* + * If we made some space in the resource, lease to as many queued guys that + * will fit now, respecting their queuing order and amount they want + */ + +void +sais_resource_check_if_can_accept_queued(sai_resource_wellknown_t *wk) +{ + assert(wk->allocated <= wk->budget); + assert(wk->allocated >= 0); + + if (wk->budget == wk->allocated) + return; + + while (wk->owner_queued.count) { + sai_resource_requisition_t *rr = + lws_container_of(wk->owner_queued.head, + sai_resource_requisition_t, + list_resource_queued_leased); + struct pss *pss = lws_container_of(rr->list_pss.owner, + struct pss, res_owner); + sai_resource_msg_t *m; + + if (wk->budget - wk->allocated < (long)rr->amount) + /* + * Have to be allocated in the queued order, so not + * being able to do this one blocks everything else + */ + return; + + wk->allocated = wk->allocated + (long)rr->amount; + + lwsl_notice("%s: leased %u to %s, -> %lu/%lu\n", __func__, + rr->amount, rr->cookie, + wk->allocated, wk->budget); + + rr->allocated_since_time = time(NULL); + lws_dll2_remove(&rr->list_resource_queued_leased); + lws_dll2_add_tail(&rr->list_resource_queued_leased, + &wk->owner_leased); + + lws_sul_schedule(wk->cx, 0, &rr->sul_expiry, sais_res_expiry, + rr->lease_secs * LWS_US_PER_SEC); + + /* + * We need to create the acceptance message + */ + + m = malloc(sizeof(*m) + LWS_PRE + 256); + if (!m) + return; + + memset(m, 0, sizeof(*m)); + + m->msg = (const char *)&m[1] + LWS_PRE; + m->len = (size_t)lws_snprintf((char *)&m[1] + LWS_PRE, 256, + "{\"schema\":\"com-warmcat-sai-resource\"," + "\"cookie\":\"%s\"," + "\"amount\":%u}", rr->cookie, rr->amount); + + lws_dll2_add_tail(&m->list, &pss->res_pending_reply_owner); + + lws_callback_on_writable(pss->wsi); + } +} + +void +sais_resource_wellknown_remove_pss(sais_t *sais, struct pss *pss) +{ + /* + * Destroy every pending resource request that belongs to this pss + */ + + lws_start_foreach_dll_safe(struct lws_dll2 *, p, p1, + pss->res_owner.head) { + sai_resource_requisition_t *rr; + + rr = lws_container_of(p, sai_resource_requisition_t, list_pss); + assert(rr->list_resource_wellknown.owner); + sais_resource_rr_destroy(rr); + + } lws_end_foreach_dll_safe(p, p1); + + /* + * Clean up any pending return resource JSON on the pss we're not + * going to get a chance to send any more + */ + + lws_start_foreach_dll_safe(struct lws_dll2 *, p, p1, + pss->res_pending_reply_owner.head) { + sai_resource_msg_t *rm; + + rm = lws_container_of(p, sai_resource_msg_t, list); + lws_dll2_remove(&rm->list); + free(rm); + + } lws_end_foreach_dll_safe(p, p1); +} + diff --git a/src/server/s-ws-builder.c b/src/server/s-ws-builder.c index b92b6f0..48bd0e2 100644 --- a/src/server/s-ws-builder.c +++ b/src/server/s-ws-builder.c @@ -1,5 +1,5 @@ /* - * Sai server - ./src/server/ws-json-rx.c + * Sai server - ./src/server/s-ws-builder.c * * Copyright (C) 2019 - 2020 Andy Green <andy@warmcat.com> * @@ -18,8 +18,8 @@ * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, * MA 02110-1301 USA * - * These are ws rx and tx handlers related to builder ws connections, on - * /builder + * These are ws rx and tx handlers related to builder ws connections, at the + * sai-server */ #include <libwebsockets.h> @@ -57,13 +57,16 @@ static const lws_struct_map_t lsm_schema_map_ba[] = { "com.warmcat.sai.taskrej"), LSM_SCHEMA (sai_artifact_t, NULL, lsm_artifact, "com-warmcat-sai-artifact"), + LSM_SCHEMA (sai_resource_t, NULL, lsm_resource, + "com-warmcat-sai-resource"), }; enum { SAIM_WSSCH_BUILDER_PLATS, SAIM_WSSCH_BUILDER_LOGS, SAIM_WSSCH_BUILDER_TASKREJ, - SAIM_WSSCH_BUILDER_ARTIFACT + SAIM_WSSCH_BUILDER_ARTIFACT, + SAIM_WSSCH_BUILDER_RESOURCE_REQ }; static void @@ -233,9 +236,12 @@ int sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t bl) { char event_uuid[33], s[128], esc[96]; + sai_resource_requisition_t *rr; + sai_resource_wellknown_t *wk; struct lwsac *ac = NULL; sai_plat_t *build, *cb; sai_rejection_t *rej; + sai_resource_t *res; lws_dll2_owner_t o; sai_artifact_t *ap; sai_task_t *task; @@ -671,6 +677,105 @@ bail: goto afail; } + break; + + case SAIM_WSSCH_BUILDER_RESOURCE_REQ: + res = (sai_resource_t *)pss->a.dest; + + /* + * We get resource requests here, and also the handing back of + * assigned leases. The requests have the resname member and + * the lease yield messages don't. + */ + + if (!res->resname) { + sai_resource_requisition_t *rr; + + /* + * An assigned resource lease is being yielded + */ + + rr = sais_resource_lookup_lease_by_cookie(&vhd->server, + res->cookie); + if (!rr) { + /* + * He never got allocated... if he's on the + * queue delete him from there... if he doesn't + * exist on our side it's OK, just finish + */ + sais_resource_destroy_queued_by_cookie( + &vhd->server, res->cookie); + + return 0; + } + + /* + * Destroy the requisition, freeing any leased resources + * allocated to him + */ + + sais_resource_rr_destroy(rr); + + return 0; + } + + /* + * This is a new request for resources, find out the well-known + * resource to attach it to + */ + + + wk = sais_resource_wellknown_by_name(&pss->vhd->server, + res->resname); + if (!wk) { + sai_resource_msg_t *mq; + + /* + * Requested well-known resource doesn't exist + */ + + lwsl_warn("%s: resource %s not well-known\n", __func__, + res->resname); + + mq = malloc(sizeof(*mq) + LWS_PRE + 256); + if (!mq) + return 0; + + memset(mq, 0, sizeof(*mq)); + + /* return with cookie but no amount == fail */ + + mq->len = (size_t)lws_snprintf((char *)&mq[1] + LWS_PRE, 256, + "{\"schema\":\"com-warmcat-sai-resource\"," + "\"cookie\":\"%s\"}", res->cookie); + + lws_dll2_add_tail(&mq->list, &pss->res_pending_reply_owner); + lws_callback_on_writable(pss->wsi); + + return 0; + } + + /* + * Create and queue the request on the right well-known + * resource manager, check if we can accept it + */ + + rr = malloc(sizeof(*rr) + strlen(res->cookie) + 1); + if (!rr) + return 0; + memset(rr, 0, sizeof(*rr)); + memcpy((char *)&rr[1], res->cookie, strlen(res->cookie) + 1); + + rr->cookie = (char *)&rr[1]; + rr->lease_secs = res->lease; + rr->amount = res->amount; + + lws_dll2_add_tail(&rr->list_pss, &pss->res_owner); + lws_dll2_add_tail(&rr->list_resource_wellknown, &wk->owner); + lws_dll2_add_tail(&rr->list_resource_queued_leased, &wk->owner_queued); + + sais_resource_check_if_can_accept_queued(wk); + break; } @@ -722,6 +827,32 @@ sais_ws_json_tx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, goto send_json; } + /* + * resource response? + */ + + if (pss->res_pending_reply_owner.count) { + sai_resource_msg_t *rm = lws_container_of(pss->res_pending_reply_owner.head, + sai_resource_msg_t, list); + + n = (int)rm->len; + if (n > lws_ptr_diff(end, p)) + n = lws_ptr_diff(end, p); + + memcpy(p, rm->msg, (unsigned int)n); + w = (size_t)n; + + lwsl_notice("%s: issuing pending resouce reply %.*s\n", __func__, (int)n, (const char *)start); + + lws_dll2_remove(&rm->list); + free(rm); + + first = 1; + pss->walk = NULL; + + goto send_json; + } + if (!pss->issue_task_owner.count) return 0; /* nothing to send */
Page fetched 0s ago, creation time: 11ms (vhost etag hits: 0%, cache hits: 0%)