Project homepage Mailing List  Warmcat.com  API Docs  Github Mirror 
    npro  
 Modern all-safe Rust Network Protocol library supporting h1, h2, h3, ws, wt sans-IO and with socket IO + tls
git clone https://npro.rs/repo/npro
 
root / src / jig / j-server.c
Author[]Andy Green <andy@warmcat.com> 2020-11-30 07:16 UTC
Committer[]Andy Green <andy@warmcat.com> 2020-11-30 09:25 UTC
Treec28605bda303893bd536020e76f8cfaa54e43f42   Raw Patch
 
clean: remove duped web stuff in server and reenable artifacts
clean: remove duped web stuff in server and reenable artifacts
diff --git a/README.md b/README.md index db359ff..08b2161 100644 --- a/README.md +++ b/README.md @@ -3,19 +3,25 @@ [![CI status](https://warmcat.com/sai/status/sai)](https://warmcat.com/git/sai) `Sai` (pronouced like 'sigh', "Trial" in Japanese) is a very lightweight -network-aware distributed CI builder and coordinating server. You can run the -sai-builder daemon on any number of devices to offer builds for that platform... -builders can run on native boxes, inside systemd-nspawn contexts, inside VMs -(eg, via qemu) for non-native arches, or on connected embedded devices which can -be flashed and run the built results, controlled by gpio and serial. - -A self-assembling constellation of Sai Builder clients make their own -connections to one or more Sai Servers, who then receive hook notifications, -read JSON from the project describing what set of build variations and tests it -should run on which platforms, and distributes work concurrently over idle -builders that have the required environment. A parallel Sai-web server is -available usually on :443 or via a proxy to provide a live web / websockets -interface with synamic updates and realtime build logs. +lws-based network-aware distributed CI builder and coordinating server. +You can run the sai-builder daemon on any number of devices ad-hoc without +central registration or inbound internet access, to offer builds for those +platforms... builders can run: + + - on native boxes, + - inside systemd-nspawn contexts, + - inside VMs (eg, via qemu) for native and non-native arches, or + - cross-build against connected embedded devices which can be flashed and run the build + results, controlled by gpio and serial. + +A sai-server daemon runs on a server to receive wss connections from the builders, +git update hooks POST signed JSON job matrices from configured git servers, and +sai-server coordinates dispatching concurrent jobs to dynamically availabe remote +builders of the correct platforms, collecting logs and results. + +A sai-web server daemon is also available usually on :443 or via a proxy to provide +a live web / websockets interface with synamic updates and realtime build logs in +the browser, with JWT-authentication for manual job control. ![sai overview](./READMEs/sai-overview.png) diff --git a/assets/sai.js b/assets/sai.js index e8c2430..78a00d6 100644 --- a/assets/sai.js +++ b/assets/sai.js @@ -1022,7 +1022,7 @@ function ws_open_sai() // if (msg.data.length < 10) // return; jso = JSON.parse(msg.data); - // console.log(jso.schema); + console.log(jso.schema); if (jso.alang) { var a = jso.alang.split(","), n; diff --git a/src/builder/b-artifacts.c b/src/builder/b-artifacts.c index 08d6634..c040ade 100644 --- a/src/builder/b-artifacts.c +++ b/src/builder/b-artifacts.c @@ -40,7 +40,7 @@ saib_artifact_rx(void *userobj, const uint8_t *buf, size_t len, int flags) { // sai_artifact_t *ap = (sai_artifact_t *)userobj; - return 0; + return LWSSSSRET_OK; } static lws_ss_state_return_t @@ -62,7 +62,7 @@ saib_artifact_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, LWS_ARRAY_SIZE(lsm_schema_json_map_artifact), 0, ap); if (!js) - return -1; + return LWSSSSRET_DESTROY_ME; lws_struct_json_serialize(js, buf, *len, &w); *len = w; @@ -71,11 +71,13 @@ saib_artifact_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, lws_ss_request_tx(ap->ss); lwsl_notice("%s: sent JSON %s\n", __func__, (const char *)buf); - return 0; + return LWSSSSRET_OK; } - if (ap->fd == -1) - return 1; /* nothing to send */ + if (ap->fd == -1) { + lwsl_info("%s: completion with fd = -1\n", __func__); + lws_ss_start_timeout(ap->ss, 5 * LWS_US_PER_SEC); + } n = read(ap->fd, buf, #if defined(WIN32) @@ -88,16 +90,16 @@ saib_artifact_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, *len = 0; close(ap->fd); ap->fd = -1; - return -1; + return LWSSSSRET_DESTROY_ME; } if (!n) { lwsl_notice("%s: file EOF\n", __func__); - return 1; /* nothing to send */ + return LWSSSSRET_TX_DONT_SEND; /* nothing to send */ } ap->ofs += n; - lwsl_debug("%s: %p: writing %d at +%llu / %llu\n", __func__, ap->ss, n, - (unsigned long long)ap->ofs, (unsigned long long)ap->len); + lwsl_info("%s: %p: writing %d at +%llu / %llu (0x%02X)\n", __func__, ap->ss, n, + (unsigned long long)ap->ofs, (unsigned long long)ap->len, buf[0]); *len = (size_t)n; if (ap->ofs == ap->len) { lwsl_notice("%s: reached logical end of artifact\n", __func__); @@ -109,7 +111,7 @@ saib_artifact_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, /* even if we finished, we want to come back to close */ lws_ss_request_tx(ap->ss); - return 0; + return LWSSSSRET_OK; } static lws_ss_state_return_t @@ -160,16 +162,20 @@ saib_artifact_state(void *userobj, void *sh, lws_ss_constate_t state, lws_ss_request_tx_len(ap->ss, (unsigned long)ap->len); break; + case LWSSSCS_TIMEOUT: + lwsl_info("%s: timeout\n", __func__); + return LWSSSSRET_DESTROY_ME; + case LWSSSCS_DISCONNECTED: // lwsl_notice("%s: LWSSSCS_DISCONNECTED\n", __func__); /* don't retry */ - return -1; + return LWSSSSRET_DESTROY_ME; default: break; } - return 0; + return LWSSSSRET_OK; } const lws_ss_info_t ssi_sai_artifact = { diff --git a/src/builder/b-sai.c b/src/builder/b-sai.c index 0d61db2..67e8144 100644 --- a/src/builder/b-sai.c +++ b/src/builder/b-sai.c @@ -101,6 +101,7 @@ static const char * const default_ss_policy = "\"http_url\":" "\"\"," /* filled in by url */ "\"tls\":" "true," "\"opportunistic\":" "true," + "\"ws_binary\":" "true," /* we're sending binary */ "\"retry\":" "\"default\"," "\"metadata\": [" "{\"url\": \"\"}" diff --git a/src/server/s-comms.c b/src/server/s-comms.c index c2937c7..fe88a57 100644 --- a/src/server/s-comms.c +++ b/src/server/s-comms.c @@ -290,17 +290,11 @@ typedef enum { SHMUT_NONE = -1, SHMUT_HOOK, SHMUT_BROWSE, - SHMUT_STATUS, - SHMUT_ARTIFACTS, - SHMUT_LOGIN } sai_http_murl_t; static const char * const well_known[] = { "/update-hook", "/sai/browse", - "/status", - "/artifacts/", /* HTTP api for accessing build artifacts */ - "/login" }; static const char *hmac_names[] = { @@ -332,26 +326,6 @@ sai_get_head_status(struct vhd *vhd, const char *projname) return state; } - -static int -sai_login_cb(void *data, const char *name, const char *filename, - char *buf, int len, enum lws_spa_fileupload_states state) -{ - return 0; -} - -static const char * const auth_param_names[] = { - "lname", - "lpass", - "success_redir", -}; - -enum enum_param_names { - EPN_LNAME, - EPN_LPASS, - EPN_SUCCESS_REDIR, -}; - static int callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, void *in, size_t len) @@ -362,9 +336,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; - char projname[64]; - int n, resp, r; - const char *cp; + int n, resp; (void)end; (void)p; @@ -460,33 +432,6 @@ callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, lwsl_notice("LWS_CALLBACK_HTTP: sees hook\n"); return 0; - case SHMUT_STATUS: - /* - * in is a string like /libwebsockets/status.svg - */ - cp = ((const char *)in) + 7; - while (*cp == '/') - cp++; - n = 0; - while (*cp != '/' && *cp && (size_t)n < sizeof(projname) - 1) - projname[n++] = *cp++; - projname[n] = '\0'; - - // lwsl_notice("%s: status %s\n", __func__, projname); - - r = sai_get_head_status(vhd, projname); - if (r < 2) - r = 2; - n = lws_snprintf(projname, sizeof(projname), - "../decal-%d.svg", r); - - if (lws_http_redirect(wsi, 307, - (unsigned char *)projname, n, - &p, end) < 0) - return -1; - - goto passthru; - default: lwsl_notice("%s: DEFAULT!!!\n", __func__); return 0; @@ -502,58 +447,12 @@ callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, goto bail; goto try_to_reuse; - - case LWS_CALLBACK_HTTP_WRITEABLE: - - lwsl_notice("%s: HTTP_WRITEABLE\n", __func__); - - if (!pss || !pss->blob_artifact) - break; - - n = lws_ptr_diff(end, start); - if ((int)(pss->artifact_length - pss->artifact_offset) < n) - n = (int)(pss->artifact_length - pss->artifact_offset); - - if (sqlite3_blob_read(pss->blob_artifact, start, n, - pss->artifact_offset)) { - lwsl_err("%s: blob read failed\n", __func__); - return -1; - } - - pss->artifact_offset += n; - - if (lws_write(wsi, start, n, - pss->artifact_offset != pss->artifact_length ? - LWS_WRITE_HTTP : LWS_WRITE_HTTP_FINAL) != n) - return -1; - - if (pss->artifact_offset != pss->artifact_length) - lws_callback_on_writable(wsi); - - break; - /* * Notifcation POSTs */ case LWS_CALLBACK_HTTP_BODY: - if (pss->login_form) { - - if (!pss->spa) { - pss->spa = lws_spa_create(wsi, auth_param_names, - LWS_ARRAY_SIZE(auth_param_names), - 1024, sai_login_cb, pss); - if (!pss->spa) { - lwsl_err("failed to create spa\n"); - return -1; - } - } - - goto spa_process; - - } - if (!pss->our_form) { lwsl_notice("%s: not our form\n", __func__); goto passthru; @@ -616,8 +515,6 @@ callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, } } -spa_process: - /* let it parse the POST data */ if (!pss->spa_failed && @@ -634,7 +531,7 @@ spa_process: lwsl_user("%s: LWS_CALLBACK_HTTP_BODY_COMPLETION: %d\n", __func__, (int)len); - if (!pss->our_form && !pss->login_form) { + if (!pss->our_form) { lwsl_user("%s: no sai form\n", __func__); goto passthru; } @@ -670,25 +567,10 @@ spa_process: return 0; /* - * ws connections from builders and browsers + * ws connections from builders */ case LWS_CALLBACK_FILTER_PROTOCOL_CONNECTION: - n = lws_hdr_copy(wsi, (char *)buf, sizeof(buf) - 1, - WSI_TOKEN_GET_URI); - if (!n) - buf[0] = '\0'; - //lwsl_notice("%s: checking with lwsgs for ws conn: %s\n", - // __func__, (const char *)buf); - - /* - * Builders don't authenticate using sessions... - */ - - if (n >= 8 && !strncmp((const char *)buf + n - 8, - "/builder", 8)) - return 0; - return 0; case LWS_CALLBACK_ESTABLISHED: @@ -696,9 +578,6 @@ spa_process: pss->vhd = vhd; if (!vhd) return -1; - pss->alang[0] = '\0'; - lws_hdr_copy(wsi, pss->alang, sizeof(pss->alang), - WSI_TOKEN_HTTP_ACCEPT_LANGUAGE); if (lws_hdr_total_length(wsi, WSI_TOKEN_GET_URI)) { if (lws_hdr_copy(wsi, (char *)start, 64, @@ -793,8 +672,8 @@ spa_process: if (!pss->announced) { /* - * Update the sai-webs about the builder removal, so they - * can update their connected browsers + * Update the sai-webs about the builder removal, so + * they can update their connected browsers */ sais_list_builders(vhd); @@ -812,10 +691,7 @@ spa_process: default: passthru: - // if (!pss || !vhd) break; - - // return vhd->gsp->callback(wsi, reason, pss->pss_gs, in, len); } return lws_callback_http_dummy(wsi, reason, user, in, len); diff --git a/src/server/s-private.h b/src/server/s-private.h index 2337450..38f707d 100644 --- a/src/server/s-private.h +++ b/src/server/s-private.h @@ -75,19 +75,6 @@ typedef struct { uint8_t nondefault; } sai_notification_t; -typedef enum { - WSS_IDLE, - WSS_PREPARE_OVERVIEW, - WSS_SEND_OVERVIEW, - WSS_PREPARE_BUILDER_SUMMARY, - WSS_SEND_BUILDER_SUMMARY, - - WSS_PREPARE_TASKINFO, - WSS_SEND_ARTIFACT_INFO, - - WSS_PREPARE_EVENTINFO, - WSS_SEND_EVENTINFO, -} ws_state; typedef struct sai_builder { sais_t c; @@ -104,14 +91,6 @@ struct pss { sai_notification_t sn; struct lws_dll2 same; /* owner: vhd.builders */ - struct lws_dll2 subs_list; - - uint64_t sub_timestamp; - char sub_task_uuid[65]; - char specific[65]; - char specific_project[96]; - char auth_user[33]; - sqlite3 *pdb_artifact; sqlite3_blob *blob_artifact; @@ -143,18 +122,14 @@ struct pss { /* notification hmac information */ char notification_sig[128]; - char alang[128]; struct lws_genhmac_ctx hmac; enum lws_genhmac_types hmac_type; char our_form; - char login_form; uint64_t first_log_timestamp; uint64_t artifact_offset; uint64_t artifact_length; - ws_state send_state; - unsigned int spa_failed:1; unsigned int subsequent:1; /* for individual JSON */ unsigned int dry:1; @@ -283,3 +258,6 @@ sais_set_task_state(struct vhd *vhd, const char *builder_name, void 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); diff --git a/src/server/s-task.c b/src/server/s-task.c index 66e26d3..f36737f 100644 --- a/src/server/s-task.c +++ b/src/server/s-task.c @@ -28,7 +28,7 @@ #include "s-private.h" -static int +int sql3_get_integer_cb(void *user, int cols, char **values, char **name) { unsigned int *pui = (unsigned int *)user; diff --git a/src/server/s-ws-builder.c b/src/server/s-ws-builder.c index 44936db..768839d 100644 --- a/src/server/s-ws-builder.c +++ b/src/server/s-ws-builder.c @@ -243,6 +243,12 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b uint64_t rid; int n, m; + if (pss->bulk_binary_data) { + lwsl_info("%s: bulk %d\n", __func__, (int)bl); + m = bl; + goto handle; + } + /* * use the schema name on the incoming JSON to decide what kind of * structure to instantiate @@ -254,9 +260,6 @@ sais_ws_json_rx_builder(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t b * - received the JSON and be handling appeneded blob data */ - if (pss->bulk_binary_data) - goto handle; - if (!pss->frag) { memset(&pss->a, 0, sizeof(pss->a)); pss->a.map_st[0] = lsm_schema_map_ba; @@ -479,6 +482,8 @@ bail: case SAIM_WSSCH_BUILDER_ARTIFACT: /* + * Builder wants to send us an artifact. + * * We get sent a JSON object immediately followed by binary * data for the artifact. * @@ -486,17 +491,19 @@ bail: * artifact table. */ - lwsl_debug("%s: SAIM_WSSCH_BUILDER_ARTIFACT\n", __func__); + lwsl_info("%s: SAIM_WSSCH_BUILDER_ARTIFACT: m = %d, bl = %d\n", __func__, m, (int)bl); if (!pss->bulk_binary_data) { + lwsl_info("%s: BUILDER_ARTIFACT: blob start, m = %d\n", __func__, m); + ap = (sai_artifact_t *)pss->a.dest; sai_task_uuid_to_event_uuid(event_uuid, ap->task_uuid); /* * Open the event-specific database object... the - * handle is closed when the stream closes for whatever + * handle is closed when the stream closes, for whatever * reason. */ @@ -619,14 +626,16 @@ bail: */ pss->bulk_binary_data = 1; pss->artifact_length = ap->len; - } else + } else { m = bl; + lwsl_info("%s: BUILDER_ARTIFACT: blob bulk\n", __func__); + } if (m) { - lwsl_notice("%s: blob write +%d, ofs %llu / %llu, len %d\n", + lwsl_info("%s: blob write +%d, ofs %llu / %llu, len %d (0x%02x)\n", __func__, (int)(bl - m), (unsigned long long)pss->artifact_offset, - (unsigned long long)pss->artifact_length, m); + (unsigned long long)pss->artifact_length, m, buf[0]); if (sqlite3_blob_write(pss->blob_artifact, (uint8_t *)buf + (bl - m), (int)m, pss->artifact_offset)) { @@ -636,10 +645,29 @@ bail: lws_set_timeout(pss->wsi, PENDING_TIMEOUT_HTTP_CONTENT, 5); pss->artifact_offset += (int)m; - } + } else + lwsl_info("%s: no m\n", __func__); + + lwsl_info("%s: ofs %d, len %d\n", __func__, (int)pss->artifact_offset, (int)pss->artifact_length); if (pss->artifact_offset == pss->artifact_length) { + int state; + lwsl_notice("%s: blob upload finished\n", __func__); + pss->bulk_binary_data = 0; + + ap = (sai_artifact_t *)pss->a.dest; + + lws_sql_purify(esc, ap->task_uuid, sizeof(esc)); + lws_snprintf(s, sizeof(s)," select state from tasks where uuid == \"%s\"", esc); + if (sqlite3_exec((sqlite3 *)pss->pdb_artifact, s, + sql3_get_integer_cb, &state, NULL) != SQLITE_OK) { + lwsl_err("%s: %s: %s: fail\n", __func__, s, + sqlite3_errmsg(pss->pdb_artifact)); + goto bail; + } + + sais_taskchange(pss->vhd->h_ss_websrv, ap->task_uuid, state); goto afail; } diff --git a/src/web/w-private.h b/src/web/w-private.h index 8e3b27d..dbc2eff 100644 --- a/src/web/w-private.h +++ b/src/web/w-private.h @@ -121,6 +121,7 @@ typedef struct saiw_scheduled { uint8_t subsequent:1; /* for individual JSON */ uint8_t ov_db_done:1; /* for individual JSON */ + uint8_t logsub:1; /* for individual JSON */ } saiw_scheduled_t; diff --git a/src/web/w-ws-browser.c b/src/web/w-ws-browser.c index 1db984a..560412a 100644 --- a/src/web/w-ws-browser.c +++ b/src/web/w-ws-browser.c @@ -324,17 +324,7 @@ saiw_pss_schedule_taskinfo(struct pss *pss, const char *task_uuid, int logsub) goto bail; } - /* does he want to subscribe to logs? */ - if (logsub) { - strcpy(pss->sub_task_uuid, sch->one_task->uuid); - lws_dll2_add_head(&pss->subs_list, &pss->vhd->subs_owner); - pss->sub_timestamp = 0; /* where we got up to */ - lws_callback_on_writable(pss->wsi); - - lwsl_notice("%s: subscribed to logs for %s\n", __func__, - pss->sub_task_uuid); - } - + sch->logsub = logsub; sch->one_event = lws_container_of(o.head, sai_event_t, list); saiw_alloc_sched(pss, WSS_PREPARE_BUILDER_SUMMARY); @@ -1033,6 +1023,8 @@ b_finish: * (all in .query_ac) */ + lwsl_info("%s: PREPARE_TASKINFO: one_task %p\n", __func__, sch->one_task); + task_reply.event = sch->one_event; task_reply.task = sch->one_task; task_reply.auth_secs = pss->authorized ? pss->expiry_unix_time - lws_now_secs() : 0; @@ -1089,8 +1081,9 @@ b_finish: lwsl_debug("%s: ---------------- no artifacts\n", __func__); /* there's no artifact stuff to do */ endo = 1; - } - sch->one_task = NULL; + } else + lwsl_debug("%s: WSS_PREPARE_TASKINFO: planning on artifacts\n", __func__); + // sch->one_task = NULL; if (n == LSJS_RESULT_ERROR) { lwsl_notice("%s: taskinfo: error generating json\n", __func__); return 1; @@ -1106,6 +1099,8 @@ b_finish: if (sch->owner.head) { sai_artifact_t *aft = (sai_artifact_t *)sch->owner.head; + lwsl_info("%s: WSS_SEND_ARTIFACT_INFO: consuming artifact\n", __func__); + lws_dll2_remove(&aft->list); /* we don't want to disclose this to browsers */ @@ -1142,7 +1137,22 @@ b_finish: send_it: flags = lws_write_ws_flags(LWS_WRITE_TEXT, first, endo || lg || !sch->walk); - if (lg || !sch->walk || endo) { + if (lg || endo || + (pss->send_state != WSS_SEND_ARTIFACT_INFO && !sch->walk) || + (pss->send_state == WSS_SEND_ARTIFACT_INFO && + (!sch || !sch->owner.head))) { + + /* does he want to subscribe to logs? */ + if (sch && sch->logsub && sch->one_task) { + strcpy(pss->sub_task_uuid, sch->one_task->uuid); + lws_dll2_add_head(&pss->subs_list, &pss->vhd->subs_owner); + pss->sub_timestamp = 0; /* where we got up to */ + lws_callback_on_writable(pss->wsi); + + lwsl_info("%s: subscribed to logs for %s\n", __func__, + pss->sub_task_uuid); + } + pss->send_state = WSS_IDLE; saiw_dealloc_sched(sch); } @@ -1150,14 +1160,6 @@ send_it: if (lws_write(pss->wsi, start, p - start, flags) < 0) return -1; - /* - * We get a bad ratio of reads to write when the builder spams us - * with rx... we have to try to clear as much as we can in one go. - */ - -// if (!lws_send_pipe_choked(pss->wsi)) -// goto again; - lws_callback_on_writable(pss->wsi); return 0;
Page fetched 0s ago, creation time: 8ms (vhost etag hits: 0%, cache hits: 0%)