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.
+
+
+
+## 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 */