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 / web / w-artifact.c
Author[]Andy Green <andy@warmcat.com> 2026-05-04 08:23 UTC
Committer[]Andy Green <andy@warmcat.com> 2026-05-05 04:11 UTC
Tree4b48c0c1f9cdd0e0e3f1ee1569db02fa72625453   Raw Patch
 
watchers
watchers
diff --git a/READMEs/README-watchers.md b/READMEs/README-watchers.md new file mode 100644 index 0000000..fc3c993 --- /dev/null +++ b/READMEs/README-watchers.md @@ -0,0 +1,114 @@ +# Sai Watcher System + +The Sai Watcher system allows the CI server to monitor external asynchronous services (like Coverity, SonarCloud, etc.) via periodic HTML screenscraping. This is useful for services that provide status updates at a public URL after a build or upload is completed. + +## How it Works + +1. **Trigger**: A builder task reports a public status URL by printing `SAI_WATCH_URL: <url>` to its log (stdout or stderr). +2. **Registration**: The server detects this prefix, identifies the service type by matching the URL against configured patterns, and creates a record in the central `watchers` table. +3. **Polling**: A background timer periodically fetches the HTML from the URL using Secure Streams. +4. **Scraping**: The server applies configured rules (prefix/suffix/anchor) to extract metrics from the HTML. +5. **UI Rendering**: The extracted metrics are sent to the browser, which renders them dynamically according to UI rules provided in the service configuration. + +## Configuration + +Watchers are defined in JSON files located in `/etc/sai/server/conf.d/`. Each file can contain a `"watchers"` array. + +### Service Object Fields + +| Field | Type | Description | +| :--- | :--- | :--- | +| `name` | string | Unique name for the service (e.g., "coverity"). Used to find icons in `assets/watchers/<name>/icon.svg`. | +| `match` | string | Substring to match against the reported URL to identify this service. | +| `rules` | array | List of scraping rules to extract data from HTML. | +| `ui` | array | List of rendering rules for the Web UI. | + +### Scraping Rule Fields (`rules`) + +| Field | Type | Description | +| :--- | :--- | :--- | +| `label` | string | Key name for the extracted value in the metrics JSON. | +| `prefix` | string | HTML string immediately preceding the value. | +| `suffix` | string | HTML string immediately following the value. | +| `anchor` | string | (Optional) HTML string to find first before looking for the prefix. | +| `final` | boolean | If this rule matches, the watcher transitions to the `FINISHED` state and stops polling. | + +### UI Rule Fields (`ui`) + +| Field | Type | Description | +| :--- | :--- | :--- | +| `label` | string | Human-readable label displayed in the UI. | +| `key` | string | The metric key (from `rules.label`) to display. | +| `warn_if_gt` | number | (Optional) Value above which the metric is highlighted as a warning. | +| `fail_if_gt` | number | (Optional) Value above which the metric is highlighted as a failure. | + +## Example: Coverity + +To configure a watcher for Coverity, create a file like `/etc/sai/server/conf.d/coverity.json`: + +```json +{ + "watchers": [ + { + "name": "coverity", + "match": "scan.coverity.com/projects/", + "rules": [ + { + "label": "status", + "prefix": "Last Build Status:</td><td>", + "suffix": "</td>" + }, + { + "label": "density", + "prefix": "Defect Density:</td><td>", + "suffix": "</td>" + }, + { + "label": "defects", + "prefix": "Outstanding Defects:</td><td>", + "suffix": "</td>" + }, + { + "label": "passed", + "prefix": "Last Build Status:</td><td>Passed", + "suffix": "</td>", + "final": true + } + ], + "ui": [ + { + "label": "Status", + "key": "status" + }, + { + "label": "Density", + "key": "density" + }, + { + "label": "Defects", + "key": "defects", + "warn_if_gt": 0, + "fail_if_gt": 10 + } + ] + } + ] +} +``` + +### Triggering from `.sai.json` + +In your task steps, after a successful upload, simply echo the status URL: + +```json +"steps": [ + { + "name": "upload", + "run": "<upload commands> ... && echo \"SAI_WATCH_URL: https://scan.coverity.com/projects/my-project\"" + } +] +``` + +## Assets + +Place a service icon at `assets/watchers/<name>/icon.svg` on the server to have it appear in the dashboard. diff --git a/assets/sai.css b/assets/sai.css index 045994a..9016e4b 100644 --- a/assets/sai.css +++ b/assets/sai.css @@ -1300,3 +1300,48 @@ canvas.power-graph { font-size: 0.8em; color: #555; } +.watchers-row { + display: flex; + flex-wrap: wrap; + padding: 4px; + gap: 8px; +} + +.watcher { + display: flex; + align-items: center; + background: rgba(255, 255, 255, 0.4); + border-radius: 4px; + padding: 2px 4px; + border: 1px solid rgba(0, 0, 0, 0.1); +} + +.watcher-icon { + width: 20px; + height: 20px; + margin-right: 4px; + vertical-align: middle; +} + +.watcher-metrics { + display: flex; + gap: 4px; + font-size: 8pt; + font-weight: bold; +} + +.watcher-metric { + background: rgba(0, 0, 0, 0.05); + padding: 1px 4px; + border-radius: 3px; +} + +.watcher-fail { + background: #ff4136; + color: white; +} + +.watcher-warn { + background: #ffdc00; + color: black; +} diff --git a/assets/sai.js b/assets/sai.js index fe786f3..787c0aa 100644 --- a/assets/sai.js +++ b/assets/sai.js @@ -397,6 +397,7 @@ var lang_zhs = "{" + var logs = "", redpend = 0, gitohashi_integ = 0, authd = 0, auth_is_admin = 0, auth_grant_level = -1, exptimer, auth_user = "", logAnsiState = {}, logs_pending = "", lines_pending = "", times_pending = "", ongoing_task_activities = {}, last_log_timestamp = 0, spreadsheet_data_cache = {}, loadreport_data_cache = {}, + watcher_services = [], fadingTasks = new Map(); var segment_stack = []; @@ -1191,6 +1192,45 @@ function summarize_build_situation(event_uuid) }; } +function sai_watcher_render(w) { + var s = "", svc = null; + + /* Find service definition */ + if (watcher_services && watcher_services.watchers) { + watcher_services.watchers.forEach(sv => { + if (sv.name === w.service_name) svc = sv; + }); + } + + s = "<div class=\"watcher\" title=\"" + san(w.service_name) + "\">"; + s += "<a href=\"" + san(w.url) + "\" target=\"_blank\">"; + s += "<img src=\"/sai/watchers/" + san(w.service_name) + "/icon.svg\" class=\"watcher-icon\">"; + s += "</a>"; + + if (w.metrics_json) { + try { + var m = JSON.parse(w.metrics_json); + if (svc && svc.ui) { + s += "<div class=\"watcher-metrics\">"; + svc.ui.forEach(u => { + if (typeof m[u.key] !== 'undefined') { + var val = m[u.key]; + var cl = ""; + if (typeof u.fail_if_gt !== 'undefined' && parseInt(val) > u.fail_if_gt) cl = " watcher-fail"; + else if (typeof u.warn_if_gt !== 'undefined' && parseInt(val) > u.warn_if_gt) cl = " watcher-warn"; + + s += "<span class=\"watcher-metric" + cl + "\" title=\"" + san(u.label) + "\">" + san(val) + "</span>"; + } + }); + s += "</div>"; + } + } catch (e) { } + } + s += "</div>"; + + return s; +} + function sai_event_summary_render(o, now_ut, reset_all_icon) { var s, q, ctn = "", wai, s1 = "", n, e = o.e; @@ -1259,7 +1299,17 @@ function sai_event_summary_render(o, now_ut, reset_all_icon) "</td></tr><tr><td class=\"nomar e6\" id=\"sumbs-" + e.uuid + "\"></td></tr>" + "</table></td>"; } - s += "</tr><tr><td class=\"nomar e6\" colspan=\"2\" id=\"sumbs-" + e.uuid +"\"></td></tr></table>"; + s += "</tr>"; + + if (o.watchers && o.watchers.length) { + s += "<tr><td class=\"nomar\" colspan=\"2\"><div class=\"watchers-row\">"; + o.watchers.forEach(w => { + s += sai_watcher_render(w); + }); + s += "</div></td></tr>"; + } + + s += "<tr><td class=\"nomar e6\" colspan=\"2\" id=\"sumbs-" + e.uuid +"\"></td></tr></table>"; return s; } @@ -2188,6 +2238,10 @@ function ws_open_sai() } break; + case "com.warmcat.sai.watcher_services": + watcher_services = jso; + break; + case "sai.warmcat.com.overview": /* * Sent with an array of e[] to start, but also diff --git a/assets/watchers/coverity/icon.svg b/assets/watchers/coverity/icon.svg new file mode 100644 index 0000000..21aabcf --- /dev/null +++ b/assets/watchers/coverity/icon.svg @@ -0,0 +1,18 @@ +<?xml version="1.0" encoding="UTF-8" standalone="no"?> +<!-- Created with Inkscape (http://www.inkscape.org/) --> + +<svg + width="19.873205mm" + height="19.870104mm" + viewBox="0 0 19.873205 19.870104" + version="1.1" + id="svg1" + xml:space="preserve" + xmlns="http://www.w3.org/2000/svg" + xmlns:svg="http://www.w3.org/2000/svg"><defs + id="defs1" /><g + id="layer1" + transform="translate(-91.619212,-178.89038)"><path + id="path4" + style="fill:#000000;fill-opacity:1;stroke:none;stroke-width:4;stroke-linecap:round;stroke-dasharray:none;stroke-opacity:1" + d="m 99.711226,178.89037 0.01809,0.005 c -4.694239,0.88598 -8.100157,5.01668 -8.110099,9.83661 -8.7e-5,3.97837 2.329924,7.58065 5.938655,9.18135 l -3.674711,-10.06657 -0.0956,-0.48111 0.529683,0.786 2.055689,3.96255 2.453597,-6.03426 0.963765,-2.45101 0.630965,-1.85001 0.22738,-0.80151 0.0605,-0.43976 -5.2e-4,-0.4036 -0.14469,-0.45216 -0.28216,-0.43822 -0.27078,-0.26717 -0.278015,-0.0811 z m 5.511294,0.5178 -0.004,0.002 0.017,0.005 c -0.005,-0.002 -0.009,-0.005 -0.0134,-0.007 z m 0.0134,0.007 c 0.002,9.3e-4 0.004,0.002 0.007,0.003 l 5.2e-4,-0.001 z m 0.007,0.003 -0.7121,1.28726 -6.699332,17.32401 c 1.183542,0.48297 2.448222,0.73123 3.724842,0.73122 5.48761,-1.9e-4 9.93614,-4.48991 9.93634,-10.02833 h -9.93686 l 9.93686,-5.2e-4 c -1.5e-4,-4.22357 -2.58773,-7.83578 -6.24975,-9.31364 z m -7.411432,18.61127 c -0.09151,-0.0372 -0.182476,-0.0758 -0.272852,-0.11576 l 0.153996,0.0734 z" /></g></svg> diff --git a/src/common/include/private.h b/src/common/include/private.h index abe9970..44e7bdc 100644 --- a/src/common/include/private.h +++ b/src/common/include/private.h @@ -65,6 +65,55 @@ enum { SAISPRF_SIGNALLED = 0x4000, }; +typedef enum { + SAIWS_QUEUED, + SAIWS_ONGOING, + SAIWS_FINISHED, + SAIWS_FAILED +} sai_watcher_state_t; + +typedef struct sai_watcher_rule { + lws_dll2_t list; + const char *label; + const char *prefix; + const char *suffix; + const char *anchor; + uint8_t final; +} sai_watcher_rule_t; + +typedef struct sai_watcher_ui_rule { + lws_dll2_t list; + const char *label; + const char *key; + int warn_if_gt; + int fail_if_gt; +} sai_watcher_ui_rule_t; + +typedef struct sai_watcher_service { + lws_dll2_t list; + const char *name; + const char *match; + const char *icon; + lws_dll2_owner_t rules_owner; /* sai_watcher_rule_t */ + lws_dll2_owner_t ui_owner; /* sai_watcher_ui_rule_t */ +} sai_watcher_service_t; + +typedef struct sai_watcher { + lws_dll2_t list; + char service_name[32]; + char event_hash[65]; + char task_hash[65]; + char url[256]; + char metrics_json[2048]; /* scraped data */ + uint64_t created; + uint64_t last_polled; + int state; /* sai_watcher_state_t */ + + /* server side only: transient ptr to service def if available */ + const sai_watcher_service_t *service; + void *vhd; +} sai_watcher_t; + typedef struct sais_sqlite_cache { lws_dll2_t list; char uuid[65]; @@ -109,6 +158,10 @@ typedef struct sai_viewer_state { unsigned int viewers; } sai_viewer_state_t; +typedef struct { + lws_dll2_owner_t watchers; +} sai_watcher_conf_t; + typedef struct sai_platform_load { lws_dll2_t list; /* Not used, for schema mapping */ char platform_name[128]; @@ -304,6 +357,8 @@ typedef struct sai_event { sai_event_state_t state; int uid; int sec; + + lws_dll2_owner_t watcher_owner; /* sai_watcher_t */ } sai_event_t; typedef struct { @@ -726,7 +781,7 @@ extern const lws_struct_map_t lsm_schema_sq3_map_artifact[1], lsm_schema_map_ta[1], lsm_schema_map_plat_simple[1], - lsm_event[11], + lsm_event[12], lsm_task[30], lsm_log[7], lsm_artifact[8], @@ -759,7 +814,14 @@ extern const lws_struct_map_t lsm_pcon_energy_report[1], lsm_schema_pcon_energy[1], lsm_pcon_control[2], - lsm_schema_pcon_control[1]; + lsm_schema_pcon_control[1], + lsm_watcher_rule[5], + lsm_watcher_ui_rule[4], + lsm_watcher_service[5], + lsm_watcher[9], + lsm_schema_sq3_map_watcher[1], + lsm_schema_json_map_watcher[1], + lsm_watcher_conf[1]; extern const lws_ss_info_t ssi_said_logproxy; extern struct lws_ss_handle *ssh[3]; diff --git a/src/common/struct-metadata.c b/src/common/struct-metadata.c index ddb98af..34e4cd9 100644 --- a/src/common/struct-metadata.c +++ b/src/common/struct-metadata.c @@ -127,6 +127,8 @@ const lws_struct_map_t lsm_event[] = { LSM_UNSIGNED (sai_event_t, state, "state"), LSM_UNSIGNED (sai_event_t, last_updated, "last_updated"), LSM_UNSIGNED (sai_event_t, sec, "sec"), + LSM_JO_LIST (sai_event_t, watcher_owner, sai_watcher_t, list, + NULL, lsm_watcher, "watchers"), }; const lws_struct_map_t lsm_schema_json_map_event[] = { @@ -466,3 +468,51 @@ const lws_struct_map_t lsm_schema_pcon_control[] = { LSM_SCHEMA(sai_pcon_control_t, NULL, lsm_pcon_control, "com.warmcat.sai.pcon_control"), }; +const lws_struct_map_t lsm_watcher_rule[] = { + LSM_STRING_PTR (sai_watcher_rule_t, label, "label"), + LSM_STRING_PTR (sai_watcher_rule_t, prefix, "prefix"), + LSM_STRING_PTR (sai_watcher_rule_t, suffix, "suffix"), + LSM_STRING_PTR (sai_watcher_rule_t, anchor, "anchor"), + LSM_UNSIGNED (sai_watcher_rule_t, final, "final"), +}; + +const lws_struct_map_t lsm_watcher_ui_rule[] = { + LSM_STRING_PTR (sai_watcher_ui_rule_t, label, "label"), + LSM_STRING_PTR (sai_watcher_ui_rule_t, key, "key"), + LSM_SIGNED (sai_watcher_ui_rule_t, warn_if_gt, "warn_if_gt"), + LSM_SIGNED (sai_watcher_ui_rule_t, fail_if_gt, "fail_if_gt"), +}; + +const lws_struct_map_t lsm_watcher_service[] = { + LSM_STRING_PTR (sai_watcher_service_t, name, "name"), + LSM_STRING_PTR (sai_watcher_service_t, match, "match"), + LSM_STRING_PTR (sai_watcher_service_t, icon, "icon"), + LSM_LIST (sai_watcher_service_t, rules_owner, sai_watcher_rule_t, list, + NULL, lsm_watcher_rule, "rules"), + LSM_LIST (sai_watcher_service_t, ui_owner, sai_watcher_ui_rule_t, list, + NULL, lsm_watcher_ui_rule, "ui"), +}; + +const lws_struct_map_t lsm_watcher[] = { + LSM_CARRAY (sai_watcher_t, service_name, "service_name"), + LSM_CARRAY (sai_watcher_t, event_hash, "event_hash"), + LSM_CARRAY (sai_watcher_t, task_hash, "task_hash"), + LSM_CARRAY (sai_watcher_t, url, "url"), + LSM_CARRAY (sai_watcher_t, metrics_json, "metrics_json"), + LSM_UNSIGNED (sai_watcher_t, created, "created"), + LSM_UNSIGNED (sai_watcher_t, last_polled, "last_polled"), + LSM_UNSIGNED (sai_watcher_t, state, "state"), +}; + +const lws_struct_map_t lsm_schema_sq3_map_watcher[] = { + LSM_SCHEMA_DLL2 (sai_watcher_t, list, NULL, lsm_watcher, "watchers"), +}; + +const lws_struct_map_t lsm_schema_json_map_watcher[] = { + LSM_SCHEMA_DLL2 (sai_watcher_t, list, NULL, lsm_watcher, "com.warmcat.sai.watchers"), +}; + +const lws_struct_map_t lsm_watcher_conf[] = { + LSM_LIST(sai_watcher_conf_t, watchers, sai_watcher_service_t, list, + NULL, lsm_watcher_service, "watchers"), +}; diff --git a/src/server/CMakeLists.txt b/src/server/CMakeLists.txt index ee66246..66982a4 100644 --- a/src/server/CMakeLists.txt +++ b/src/server/CMakeLists.txt @@ -16,6 +16,7 @@ set(SRCS s-ws-web.c s-webops.c s-resource.c + s-watcher.c ../common/c-utils.c ../common/c-sqlite3.c ../common/struct-metadata.c diff --git a/src/server/s-central.c b/src/server/s-central.c index 3f4f583..82fd9ab 100644 --- a/src/server/s-central.c +++ b/src/server/s-central.c @@ -182,6 +182,19 @@ sais_ensure_tables(struct vhd *vhd) " pcon_name varchar(64)," " builder_name varchar(64)" ");", "create pcon_builders table"); + + sai_sqlite3_statement(vhd->server.pdb, + "CREATE TABLE IF NOT EXISTS watchers (" + " service_name varchar(32)," + " event_hash varchar(65)," + " task_hash varchar(65)," + " url varchar(256)," + " state integer," + " created integer," + " last_polled integer," + " metrics_json text," + " PRIMARY KEY (event_hash, service_name)" + ");", "create watchers table"); } @@ -293,6 +306,10 @@ sais_central_cb(lws_sorted_usec_list_t *sul) if (!vhd->sul_gc_events.list.owner) lws_sul_schedule(context, 0, &vhd->sul_gc_events, sais_central_gc_deleted_events_cb, 10 * LWS_US_PER_MS); + + if (!vhd->sul_watcher.list.owner) + lws_sul_schedule(context, 0, &vhd->sul_watcher, + sais_watcher_cb, 100 * LWS_US_PER_MS); /* check again in 1s */ diff --git a/src/server/s-comms.c b/src/server/s-comms.c index 1603681..8ee231c 100644 --- a/src/server/s-comms.c +++ b/src/server/s-comms.c @@ -128,6 +128,12 @@ s_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, return -1; } + { + const char *conf_dir = "/etc/sai/server"; + lws_pvo_get_str(in, "config-dir", &conf_dir); + sais_config_watchers(vhd, conf_dir); + } + /* * Create the listed well-known resources to be managed by the * sai-server for the builders @@ -543,3 +549,94 @@ const struct lws_protocols protocol_ws = { .per_session_data_size = sizeof(struct pss), .rx_buffer_size = 0, }; + +static int +sais_config_watchers_cb(const char *dirpath, void *opaque, struct lws_dir_entry *lde) +{ + struct vhd *vhd = (struct vhd *)opaque; + char path[256], *buf; + sai_watcher_conf_t *wc; + struct lwsac *ac = NULL; + lws_dll2_owner_t o; + int n, fd; + struct stat st; + struct lejp_ctx ctx; + lws_struct_args_t args; + + memset(&o, 0, sizeof(o)); + + if (lde->type != LDOT_FILE || !strstr(lde->name, ".json")) + return 0; + + lws_snprintf(path, sizeof(path), "%s/%s", dirpath, lde->name); + lwsl_notice("%s: parsing %s for watchers\n", __func__, path); + + fd = open(path, O_RDONLY); + if (fd < 0) + return 0; + + if (fstat(fd, &st)) { + close(fd); + return 0; + } + + buf = malloc((size_t)st.st_size); + if (!buf) { + close(fd); + return 0; + } + + if (read(fd, buf, (size_t)st.st_size) != (ssize_t)st.st_size) { + free(buf); + close(fd); + return 0; + } + close(fd); + + memset(&args, 0, sizeof(args)); + args.map_st[0] = lsm_watcher_conf; + args.map_entries_st[0] = LWS_ARRAY_SIZE(lsm_watcher_conf); + args.dest = &o; + args.dest_len = sizeof(o); + + lws_struct_json_init_parse(&ctx, lws_struct_default_lejp_cb, &args); + n = lejp_parse(&ctx, (const uint8_t *)buf, (int)st.st_size); + lejp_destruct(&ctx); + ac = args.ac; + + free(buf); + + if (n < 0 || !o.head) { + lwsac_free(&ac); + return 0; + } + + wc = lws_container_of(o.head, sai_watcher_conf_t, watchers); + + /* Move the parsed watcher services to our vhd list */ + lws_start_foreach_dll_safe(struct lws_dll2 *, p, p1, wc->watchers.head) { + sai_watcher_service_t *ws = lws_container_of(p, + sai_watcher_service_t, list); + + lwsl_notice("%s: added watcher service '%s'\n", __func__, + ws->name); + lws_dll2_remove(&ws->list); + lws_dll2_add_tail(&ws->list, &vhd->watcher_services); + + } lws_end_foreach_dll_safe(p, p1); + + lwsac_free(&ac); + + return 0; +} + +int +sais_config_watchers(struct vhd *vhd, const char *config_dir) +{ + char path[256]; + + lws_snprintf(path, sizeof(path), "%s/conf.d", config_dir); + lwsl_notice("%s: scanning %s\n", __func__, path); + + return lws_dir(path, vhd, sais_config_watchers_cb); +} diff --git a/src/server/s-private.h b/src/server/s-private.h index 0f31e8a..9f6c62d 100644 --- a/src/server/s-private.h +++ b/src/server/s-private.h @@ -232,6 +232,9 @@ struct vhd { lws_sorted_usec_list_t sul_central; /* background task allocation sul */ lws_sorted_usec_list_t sul_activity; /* activity broadcast sul */ lws_sorted_usec_list_t sul_gc_events; /* incremental GC of deleted events */ + lws_sorted_usec_list_t sul_watcher; /* generic async service watcher sul */ + + lws_dll2_owner_t watcher_services; /* sai_watcher_service_t from config */ lws_usec_t last_check_abandoned_tasks; @@ -300,6 +303,12 @@ sai_db_result_t sais_task_clear_build_and_logs(struct vhd *vhd, const char *task_uuid, int from_rejection); sai_db_result_t sais_task_rebuild_last_step(struct vhd *vhd, const char *task_uuid); + +void +sais_watcher_cb(lws_sorted_usec_list_t *sul); + +int +sais_config_watchers(struct vhd *vhd, const char *config_dir); int sais_task_cancel(struct vhd *vhd, const char *task_uuid, int erase); diff --git a/src/server/s-watcher.c b/src/server/s-watcher.c new file mode 100644 index 0000000..3a21f18 --- /dev/null +++ b/src/server/s-watcher.c @@ -0,0 +1,236 @@ +#include <libwebsockets.h> +#include <string.h> +#include <time.h> +#include "s-private.h" + +typedef struct watcher_fetch { + struct lws_ss_handle *ss; + struct vhd *vhd; + sai_watcher_t *watcher; + struct lws_buflist *bl_rx; + int status; +} watcher_fetch_t; + +static lws_ss_state_return_t +sais_watcher_ss_rx(void *userobj, const uint8_t *buf, size_t len, int flags) +{ + watcher_fetch_t *f = (watcher_fetch_t *)userobj; + + if (lws_buflist_append_segment(&f->bl_rx, buf, len)) + return LWSSSSRET_DESTROY_ME; + + return LWSSSSRET_OK; +} + +static lws_ss_state_return_t +sais_watcher_ss_tx(void *userobj, lws_ss_tx_ordinal_t ord, uint8_t *buf, + size_t *len, int *flags) +{ + *len = 0; + return LWSSSSRET_OK; +} + +static void +sais_watcher_scrape(watcher_fetch_t *f) +{ + sai_watcher_t *w = f->watcher; + const sai_watcher_service_t *s = w->service; + const uint8_t *p; + size_t len; + char metrics[2048], *m = metrics, *mend = metrics + sizeof(metrics) - 1; + int first = 1; + + if (!s) + return; + + lwsl_notice("%s: scraping %s for %s\n", __func__, w->url, s->name); + + lws_snprintf(metrics, sizeof(metrics), "{"); + m = metrics + 1; + + lws_start_foreach_dll(struct lws_dll2 *, pr, s->rules_owner.head) { + sai_watcher_rule_t *r = lws_container_of(pr, sai_watcher_rule_t, list); + const char *val = NULL; + char valbuf[256]; + + /* + * This is a very simple scraper. It looks for prefix and suffix. + * If an anchor is provided, it first finds the anchor. + */ + p = NULL; + len = lws_buflist_next_segment_len(&f->bl_rx, (uint8_t **)&p); + while (p) { + const char *found = NULL; + const char *sp = (const char *)p; + + if (r->anchor) { + const char *a = strstr(sp, r->anchor); + if (a) { + /* Found anchor, now look for prefix near it */ + /* For now, just look after it. In some cases we might need to look before. */ + found = strstr(a, r->prefix); + } + } else { + found = strstr(sp, r->prefix); + } + + if (found) { + const char *start = found + strlen(r->prefix); + const char *end = strstr(start, r->suffix); + + if (end) { + size_t vlen = (size_t)lws_ptr_diff(end, start); + if (vlen >= sizeof(valbuf)) + vlen = sizeof(valbuf) - 1; + memcpy(valbuf, start, vlen); + valbuf[vlen] = '\0'; + val = valbuf; + break; + } + } + + lws_buflist_use_segment(&f->bl_rx, len); + len = lws_buflist_next_segment_len(&f->bl_rx, (uint8_t **)&p); + } + + if (val) { + m += lws_snprintf(m, lws_ptr_diff_size_t(mend, m), + "%s\"%s\":\"%s\"", first ? "" : ",", + r->label, val); + first = 0; + + if (r->final) + w->state = SAIWS_FINISHED; + } + + } lws_end_foreach_dll(pr); + + lws_snprintf(m, lws_ptr_diff_size_t(mend, m), "}"); + + lwsl_notice("%s: metrics: %s\n", __func__, metrics); + lws_strncpy(w->metrics_json, metrics, sizeof(w->metrics_json)); + + if (w->state == SAIWS_QUEUED) + w->state = SAIWS_ONGOING; + + { + struct vhd *vhd = f->vhd; + char q[2048 + 256], esc_metrics[2048 + 128]; + lws_sql_purify(esc_metrics, w->metrics_json, sizeof(esc_metrics)); + + lws_snprintf(q, sizeof(q), + "UPDATE watchers SET metrics_json='%s', last_polled=%llu, state=%d WHERE url='%s'", + esc_metrics, (unsigned long long)w->last_polled, w->state, w->url); + + if (sai_sqlite3_statement(vhd->server.pdb, q, "update watcher metrics")) + lwsl_err("%s: failed to update watcher metrics\n", __func__); + } +} + +static lws_ss_state_return_t +sais_watcher_ss_state(void *userobj, void *sh, lws_ss_constate_t state, + lws_ss_tx_ordinal_t ack) +{ + watcher_fetch_t *f = (watcher_fetch_t *)userobj; + + switch (state) { + case LWSSSCS_CREATING: + f->vhd = (struct vhd *)f->watcher->vhd; + return lws_ss_client_connect(f->ss); + + case LWSSSCS_CONNECTED: + break; + + case LWSSSCS_DISCONNECTED: + if (f->bl_rx) { + sais_watcher_scrape(f); + lws_buflist_destroy_all_segments(&f->bl_rx); + } + /* Update DB here? */ + break; + + case LWSSSCS_ALL_RETRIES_FAILED: + case LWSSSCS_DESTROYING: + lws_buflist_destroy_all_segments(&f->bl_rx); + break; + + default: + break; + } + + return LWSSSSRET_OK; +} + +const lws_ss_info_t ssi_watcher = { + .handle_offset = offsetof(watcher_fetch_t, ss), + .opaque_user_data_offset = offsetof(watcher_fetch_t, watcher), + .streamtype = "watcher", + .rx = sais_watcher_ss_rx, + .tx = sais_watcher_ss_tx, + .state = sais_watcher_ss_state, + .user_alloc = sizeof(watcher_fetch_t), +}; + +void +sais_watcher_cb(lws_sorted_usec_list_t *sul) +{ + struct vhd *vhd = lws_container_of(sul, struct vhd, sul_watcher); + struct lwsac *ac = NULL; + lws_dll2_owner_t o; + int n; + + lws_dll2_owner_clear(&o); + + /* + * Find watchers that need polling. + * We poll every 5m, 30m, or 2h depending on state. + */ + n = lws_struct_sq3_deserialize(vhd->server.pdb, + " and (state != 2)", /* not finished */ + "last_polled", + lsm_schema_sq3_map_watcher, &o, &ac, 0, 1); + + if (n > 0 && o.head) { + sai_watcher_t *w = lws_container_of(o.head, sai_watcher_t, list); + lws_usec_t next_poll = 0; + + w->vhd = vhd; + + switch (w->state) { + case SAIWS_QUEUED: + next_poll = 30 * 60 * LWS_USEC_PER_SEC; + break; + case SAIWS_ONGOING: + next_poll = 5 * 60 * LWS_USEC_PER_SEC; + break; + case SAIWS_FAILED: + next_poll = 2 * 60 * 60 * LWS_USEC_PER_SEC; + break; + } + + if (lws_now_usecs() - (lws_usec_t)w->last_polled > next_poll) { + /* Identify service */ + lws_start_foreach_dll(struct lws_dll2 *, p, vhd->watcher_services.head) { + sai_watcher_service_t *s = lws_container_of(p, sai_watcher_service_t, list); + if (strstr(w->url, s->match)) { + w->service = s; + break; + } + } lws_end_foreach_dll(p); + + if (w->service) { + lwsl_notice("%s: starting poll for %s\n", __func__, w->url); + if (lws_ss_create(vhd->context, 0, &ssi_watcher, w, NULL, NULL, NULL)) + lwsl_err("%s: failed to create ss for watcher\n", __func__); + + /* Update last_polled to avoid multiple simultaneous polls */ + w->last_polled = (uint64_t)lws_now_usecs(); + /* Update DB ... */ + } + } + } + + lwsac_free(&ac); + + lws_sul_schedule(vhd->context, 0, &vhd->sul_watcher, sais_watcher_cb, 10 * LWS_US_PER_SEC); +} diff --git a/src/server/s-ws-builder.c b/src/server/s-ws-builder.c index 95bbd54..9d957b0 100644 --- a/src/server/s-ws-builder.c +++ b/src/server/s-ws-builder.c @@ -258,6 +258,48 @@ sais_log_to_db(struct vhd *vhd, sai_log_t *log) lwsl_err("%s: failed to update build_step\n", __func__); sai_event_db_close(&vhd->sqlite3_cache, &pdb); + + if (log->len >= 14 && !memcmp(log->log, "SAI_WATCH_URL:", 14)) { + const char *url = log->log + 14; + sai_watcher_t w; + + while (*url == ' ') + url++; + + memset(&w, 0, sizeof(w)); + lws_strncpy(w.url, url, sizeof(w.url)); + lws_strncpy(w.task_hash, log->task_uuid, sizeof(w.task_hash)); + sai_task_uuid_to_event_uuid(w.event_hash, log->task_uuid); + w.created = (uint64_t)lws_now_usecs(); + w.state = SAIWS_QUEUED; + + /* Identify service immediately to store service_name */ + lws_start_foreach_dll(struct lws_dll2 *, p, vhd->watcher_services.head) { + sai_watcher_service_t *s = lws_container_of(p, sai_watcher_service_t, list); + if (strstr(w.url, s->match)) { + lws_strncpy(w.service_name, s->name, sizeof(w.service_name)); + break; + } + } lws_end_foreach_dll(p); + + if (w.service_name[0]) { + lwsl_notice("%s: triggering watcher for %s (%s)\n", __func__, w.service_name, w.url); + /* For now, just use sqlite3_exec */ + char q[512], esc_url[256], esc_svc[64], esc_event[65], esc_task[65]; + lws_sql_purify(esc_url, w.url, sizeof(esc_url)); + lws_sql_purify(esc_svc, w.service_name, sizeof(esc_svc)); + lws_sql_purify(esc_event, w.event_hash, sizeof(esc_event)); + lws_sql_purify(esc_task, w.task_hash, sizeof(esc_task)); + + lws_snprintf(q, sizeof(q), + "REPLACE INTO watchers (service_name, event_hash, task_hash, url, state, created, last_polled, metrics_json) " + "VALUES ('%s', '%s', '%s', '%s', %d, %llu, 0, '{}')", + esc_svc, esc_event, esc_task, esc_url, w.state, (unsigned long long)w.created); + + if (sai_sqlite3_statement(vhd->server.pdb, q, "insert watcher")) + lwsl_err("%s: failed to insert watcher\n", __func__); + } + } } sai_plat_t * diff --git a/src/web/w-comms.c b/src/web/w-comms.c index 7cde616..9da21d6 100644 --- a/src/web/w-comms.c +++ b/src/web/w-comms.c @@ -158,6 +158,7 @@ w_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, return -1; lwsl_err("web-callback-ws: LWS_CALLBACK_PROTOCOL_INIT\n"); + lws_dll2_owner_clear(&vhd->watcher_services); vhd->context = lws_get_context(wsi); vhd->vhost = lws_get_vhost(wsi); diff --git a/src/web/w-private.h b/src/web/w-private.h index d4d16a0..8e4d74c 100644 --- a/src/web/w-private.h +++ b/src/web/w-private.h @@ -142,6 +142,8 @@ struct vhd { lws_dll2_owner_t subs_owner; sqlite3 *pdb; + lws_dll2_owner_t watcher_services; + lws_dll2_owner_t pcon_watts_owner; unsigned int power_history[150]; int power_history_count; diff --git a/src/web/w-ws-browser.c b/src/web/w-ws-browser.c index 053a365..01188f5 100644 --- a/src/web/w-ws-browser.c +++ b/src/web/w-ws-browser.c @@ -93,6 +93,8 @@ static const lws_struct_map_t lsm_schema_json_map_bwsrx[] = { "com.warmcat.sai.stay"), LSM_SCHEMA (sai_pcon_control_t, NULL, lsm_pcon_control, /* shares struct */ "com.warmcat.sai.pcon_control"), + LSM_SCHEMA_DLL2 (sai_watcher_service_t, list, NULL, lsm_watcher_service, + "com.warmcat.sai.watcher_services"), }; enum { @@ -610,6 +612,15 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, pss->overview_offset = ti->offset; saiw_browser_broadcast_queue_builders(pss->vhd, pss); + + { + uint8_t buf[LWS_PRE + 4096], *start = buf + LWS_PRE, *p = start, *end = buf + sizeof(buf); + + p += lws_snprintf((char *)p, lws_ptr_diff_size_t(end, p), + "{\"schema\":\"com.warmcat.sai.watcher_services\",\"watchers\":[]}"); + saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start, lws_ptr_diff_size_t(p, start), LWS_WRITE_TEXT); + } + saiw_browser_queue_overview(pss->vhd, pss); break; } @@ -979,6 +990,14 @@ saiw_browser_queue_overview(struct vhd *vhd, struct pss *pss) } } + { + char wfilt[128]; + struct lwsac *ac_watchers = NULL; + lws_dll2_owner_clear(&e->watcher_owner); + lws_snprintf(wfilt, sizeof(wfilt), " and event_hash='%s'", e->uuid); + lws_struct_sq3_deserialize(vhd->pdb, wfilt, "created", + lsm_schema_sq3_map_watcher, &e->watcher_owner, &ac_watchers, 0, 0); + js = lws_struct_json_serialize_create( lsm_schema_json_map_event, LWS_ARRAY_SIZE(lsm_schema_json_map_event), 0, e); @@ -1023,6 +1042,10 @@ saiw_browser_queue_overview(struct vhd *vhd, struct pss *pss) } } while (n == LSJS_RESULT_CONTINUE); + if (ac_watchers) + lwsac_free(&ac_watchers); + } + if (lws_ptr_diff_size_t(end, p) < 128) { saiw_ws_browser_queue_REQUIRES_LWS_PRE(pss, start, lws_ptr_diff_size_t(p, start),
Page fetched 0s ago, creation time: 39ms (vhost etag hits: 0%, cache hits: 0%)