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 / assets / arch-xl6-esp32.svg
Author[]Andy Green <andy@warmcat.com> 2026-09-28 15:27 UTC
Committer[]Andy Green <andy@warmcat.com> 2026-09-28 15:29 UTC
Tree4fdc1f671e44d096fa8bd0bc4826e27c2eb22f89   Raw Patch
 
web: rss feed as JSON too, and long poll on it for changes
web: rss feed as JSON too, and long poll on it for changes

sai-web also serves the feed as JSON on [/sai]/rss.json, for tools that
want the same information without parsing RSS.  It is a sai_feed_t with
its sai_feed_item_t list, serialized with a new lws_struct map in
common, so a tool built with sai can parse it with the same map.

Both forms can now also be scoped to one repository with ?fetchurl=,
matching the fetch url its notifications give.

The feed now carries an index token, a hash of each listed event's
uuid and state name, so it changes when an event joins or leaves the
feed or changes state, but not when only task counts move.  A request
adding ?wait=<secs>&index=<token> is answered at once if the token is
out of date, and otherwise held until it is, or until the wait (capped
at 600s) expires, and then answered with the then-current feed.  So a
client can fetch the feed once, then keep asking to wait on the last
index, and only hear back when something happened.

Held requests are re-checked when sai-server reports an event change.
The check only reads the events table, so it opens no per-event dbs;
the task counts are filled in only when a feed is actually sent.

The lws reverse proxy in front of sai-web gives up on an onward request
whose response headers don't come within the context timeout_secs, so
a held request is answered with its headers at once, with no content
length (close-delimited on h1), marked as a long-lived stream with
lws_http_mark_sse(), and sent a newline every 10s until the feed
follows.  Whitespace is valid before the JSON value, and between the
XML declaration, sent up front, and <rss>.  At most 64 requests are
held at once, past that they are told to retry later with a 503.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
diff --git a/README.md b/README.md index ac8ab51..9f527c8 100644 --- a/README.md +++ b/README.md @@ -186,7 +186,9 @@ daemons (sai-server first, sai-web reconnects by itself). sai-web serves a public RSS 2.0 feed of the latest 10 events at `/sai/rss.xml`, newest notification first. The web UI links to it from the -logo area and advertises it for feed autodiscovery. It can be scoped with `?project=<name>` and / or +logo area and advertises it for feed autodiscovery. The same feed is served +as JSON at `/sai/rss.json`, for tools. Either can be scoped with any of +`?project=<name>`, `?fetchurl=<repository fetch url>` and `?branch=<branch name or full ref>`, eg, ``` @@ -218,6 +220,30 @@ Links in the feed are relative (eg, `index.html?event=<uuid>`), so they resolve against whatever url the feed was fetched from; no conf is needed to tell sai-web its public url. +### Waiting for changes (long poll) + +The feed carries an index token (`sai:index` in the channel, `"index"` in the +JSON), which changes when an event joins or leaves the feed, or any event's +state changes. Task counts changing alone don't change it. + +A request with `?wait=<secs>&index=<token>` added is answered straight away +if the token is already out of date. Otherwise it is held until the token +goes out of date, or the wait (at most 600s) runs out, and then answered with +the feed as it is then. So instead of polling, a client can fetch the feed +once and then loop asking to wait on the index from the last response, eg, + +``` +https://mydomain.com/sai/rss.json?project=libwebsockets&wait=600&index=<index from last time> +``` + +A held request gets its response headers at once, with no content length, +then a newline every 10s until the feed follows, so that proxies and +idle-connection reapers along the way leave it alone. That whitespace is +allowed before the JSON value, and between the XML declaration (which is sent +at the start) and `<rss>`, so the body is still a valid document. At most 64 +requests are held at once; past that, a request to wait is answered with a +503 and `retry-after`. + ## Build flow and support for embedded ![build flow](READMEs/sai-build-test-flow.png) diff --git a/src/common/include/private.h b/src/common/include/private.h index d20aca5..6315e8b 100644 --- a/src/common/include/private.h +++ b/src/common/include/private.h @@ -486,6 +486,45 @@ typedef struct sai_event { lws_dll2_owner_t watcher_owner; /* sai_watcher_t */ } sai_event_t; +/* + * One event as sai-web's public feed of recent events reports it, and the + * feed itself, so the JSON form of the feed (rss.json) is written by sai-web + * and read by sai-push with the same lws_struct map + */ + +typedef struct sai_feed_item { + struct lws_dll2 list; + char uuid[65]; + char project[65]; + /* the ref less any refs/heads/ */ + char branch[65]; + char ref[65]; + char hash[65]; + char fetchurl[96]; + char weburl[128]; + /* eg, "building", "succeeded", see w-rss.c */ + char state_name[16]; + char summary[96]; + /* unix time the notification creating the event arrived */ + uint64_t received; + int state; + int adhoc; + unsigned int tasks_total; + unsigned int tasks_ok; + unsigned int tasks_bad; + unsigned int tasks_building; + unsigned int tasks_wait; +} sai_feed_item_t; + +typedef struct sai_feed { + /* + * Changes whenever an event joins or leaves the feed, or an event's + * state_name changes (but not when only task counts change) + */ + char index[33]; + lws_dll2_owner_t items; /* sai_feed_item_t */ +} sai_feed_t; + typedef struct { struct lws_dll2 list; char task_uuid[65]; @@ -987,6 +1026,9 @@ extern const lws_struct_map_t lsm_schema_json_map_can[1], lsm_schema_json_map_task[1], lsm_schema_json_map_event[1], + lsm_feed_item[17], + lsm_feed[2], + lsm_schema_json_map_feed[1], lsm_resource[4], lsm_power_state[3], lsm_openshell[2], diff --git a/src/common/struct-metadata.c b/src/common/struct-metadata.c index 3d6d73d..7cf34d4 100644 --- a/src/common/struct-metadata.c +++ b/src/common/struct-metadata.c @@ -151,6 +151,36 @@ const lws_struct_map_t lsm_schema_json_map_event[] = { "com.warmcat.sai.events"), }; +const lws_struct_map_t lsm_feed_item[] = { + LSM_CARRAY (sai_feed_item_t, uuid, "uuid"), + LSM_CARRAY (sai_feed_item_t, project, "project"), + LSM_CARRAY (sai_feed_item_t, branch, "branch"), + LSM_CARRAY (sai_feed_item_t, ref, "ref"), + LSM_CARRAY (sai_feed_item_t, hash, "hash"), + LSM_CARRAY (sai_feed_item_t, fetchurl, "fetchurl"), + LSM_CARRAY (sai_feed_item_t, weburl, "weburl"), + LSM_CARRAY (sai_feed_item_t, state_name, "state_name"), + LSM_CARRAY (sai_feed_item_t, summary, "summary"), + LSM_UNSIGNED (sai_feed_item_t, received, "received"), + LSM_SIGNED (sai_feed_item_t, state, "state"), + LSM_SIGNED (sai_feed_item_t, adhoc, "adhoc"), + LSM_UNSIGNED (sai_feed_item_t, tasks_total, "tasks_total"), + LSM_UNSIGNED (sai_feed_item_t, tasks_ok, "tasks_ok"), + LSM_UNSIGNED (sai_feed_item_t, tasks_bad, "tasks_bad"), + LSM_UNSIGNED (sai_feed_item_t, tasks_building, "tasks_building"), + LSM_UNSIGNED (sai_feed_item_t, tasks_wait, "tasks_wait"), +}; + +const lws_struct_map_t lsm_feed[] = { + LSM_CARRAY (sai_feed_t, index, "index"), + LSM_LIST (sai_feed_t, items, sai_feed_item_t, list, + NULL, lsm_feed_item, "items"), +}; + +const lws_struct_map_t lsm_schema_json_map_feed[] = { + LSM_SCHEMA (sai_feed_t, NULL, lsm_feed, "com.warmcat.sai.feed"), +}; + const lws_struct_map_t lsm_schema_sq3_map_event[] = { LSM_SCHEMA_DLL2 (sai_event_t, list, NULL, lsm_event, "events"), }; diff --git a/src/web/w-comms.c b/src/web/w-comms.c index 65be1df..ade95d6 100644 --- a/src/web/w-comms.c +++ b/src/web/w-comms.c @@ -55,7 +55,9 @@ typedef enum { SHMUT_ARTIFACTS_SAI, SHMUT_LOGIN, SHMUT_RSS, - SHMUT_RSS_SAI + SHMUT_RSS_SAI, + SHMUT_RSS_JSON, + SHMUT_RSS_JSON_SAI } sai_http_murl_t; static const char * const well_known[] = { @@ -66,7 +68,9 @@ static const char * const well_known[] = { "/sai/artifacts/", /* same, via the /sai mount the pages live under */ "/login", "/rss.xml", /* public RSS 2.0 feed of recent events, see w-rss.c */ - "/sai/rss.xml" + "/sai/rss.xml", + "/rss.json", /* the same feed as JSON */ + "/sai/rss.json" }; /* @@ -513,11 +517,16 @@ w_callback_ws(struct lws *wsi, enum lws_callback_reasons reason, void *user, case SHMUT_RSS: case SHMUT_RSS_SAI: - r = saiw_rss_http(vhd, pss, wsi); + case SHMUT_RSS_JSON: + case SHMUT_RSS_JSON_SAI: + r = saiw_rss_http(vhd, pss, wsi, mu == SHMUT_RSS_JSON || + mu == SHMUT_RSS_JSON_SAI); if (r < 0) goto bail; if (!r) return 0; + if (r == 1) + goto try_to_reuse; resp = r; goto http_resp; @@ -538,7 +547,7 @@ http_resp: case LWS_CALLBACK_HTTP_WRITEABLE: - if (pss && pss->rss_tx) + if (pss && pss->rss_state) return saiw_rss_writeable(pss, wsi); if (!pss || !pss->blob_artifact) diff --git a/src/web/w-private.h b/src/web/w-private.h index a2cffb5..89871cb 100644 --- a/src/web/w-private.h +++ b/src/web/w-private.h @@ -67,6 +67,14 @@ enum { SAIM_SPECIFIC_TASK, }; +/* where a feed http transaction is up to, see w-rss.c */ +enum { + SAIW_RSS_IDLE, + SAIW_RSS_PARKED, /* held for a long poll, sending keepalives */ + SAIW_RSS_FINISHING, /* the whole response is in rss_tx */ + SAIW_RSS_FAILED, /* drop the connection */ +}; + typedef enum { SAI_AUTH_STATE_NOT_LOGGED_IN, SAI_AUTH_STATE_LOGGED_IN_NO_GRANT, @@ -136,8 +144,23 @@ struct pss { struct vhd *vhd; sqlite3 *pdb_artifact; sqlite3_blob *blob_artifact; - /* rendered rss feed waiting to go out on this http transaction */ + /* + * Feed (rss.xml / rss.json) http transaction, see w-rss.c: the + * rendered response waiting to go out, and for a held long poll + * request, its place in vhd->rss_waiters, its keepalive / deadline + * timer, and the scope and index it is waiting on + */ struct lws_buflist *rss_tx; + lws_dll2_t rss_list; + lws_sorted_usec_list_t sul_rss; + lws_usec_t rss_deadline; + char rss_project[65]; + char rss_branch[65]; + char rss_fetchurl[96]; + char rss_index[33]; + uint8_t rss_state; /* SAIW_RSS_* */ + uint8_t rss_json:1; + uint8_t rss_ka:1; lws_dll2_owner_t logs_owner; lws_sorted_usec_list_t sul_logcache; @@ -219,6 +242,7 @@ struct vhd { const char *sockpath; /* sai-server control link uds */ lws_dll2_owner_t sqlite3_cache; /* sais_sqlite_cache_t */ + lws_dll2_owner_t rss_waiters; /* pss held on long poll */ lws_dll2_owner_t tasklog_cache; }; @@ -331,7 +355,10 @@ saiw_event_summary_string(sqlite3 *pdb_event, const char *event_uuid, unsigned int *p_total); int -saiw_rss_http(struct vhd *vhd, struct pss *pss, struct lws *wsi); +saiw_rss_http(struct vhd *vhd, struct pss *pss, struct lws *wsi, int json); + +void +saiw_rss_event_change(struct vhd *vhd); int saiw_rss_writeable(struct pss *pss, struct lws *wsi); diff --git a/src/web/w-rss.c b/src/web/w-rss.c index 510bbc5..003d1d0 100644 --- a/src/web/w-rss.c +++ b/src/web/w-rss.c @@ -1,5 +1,5 @@ /* - * Sai web - RSS 2.0 feed of recent events + * Sai web - RSS 2.0 (and JSON) feed of recent events * * Copyright (C) 2019 - 2026 Andy Green <andy@warmcat.com> * @@ -18,8 +18,10 @@ * Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, * MA 02110-1301 USA * - * Public, unauthenticated RSS 2.0 feed of the most recent events, on - * [/sai]/rss.xml, optionally scoped with ?project=<repo name> and / or + * Public, unauthenticated feed of the most recent events, as RSS 2.0 on + * [/sai]/rss.xml, or the same content as JSON on [/sai]/rss.json (for tools + * like sai-push; the JSON is sai_feed_t via lws_struct). Either can be + * scoped with any of ?project=<repo name>, ?fetchurl=<repo fetch url> and * ?branch=<branch name or full ref>. * * One item per event, newest notification first. The item reflects the @@ -37,6 +39,26 @@ * Links in the feed are relative, so they resolve against the url the feed * was fetched from: whoever fetched it already knows where sai's web UI is * published, which sai-web behind the front-end proxy can't know for sure. + * + * Long poll + * --------- + * + * The feed carries an "index" token, which changes when an event joins or + * leaves the feed, or an event's state name changes (the same thing that + * changes an item guid; task counts changing alone do not change it). A + * request with ?wait=<secs>&index=<token> is answered at once if the token is + * already out of date, and otherwise held until it goes out of date or the + * wait (capped at SAIW_RSS_MAX_WAIT_S) expires, when it is answered with the + * then-current feed either way. So a client fetches the feed once, then + * loops asking to wait on the index from the last response, and only hears + * back when there is something new. + * + * Front-end proxies give up on a request whose response headers don't arrive + * promptly, and idle connections get reaped, so a held request is answered + * with its headers straight away, with no content-length, and a newline every + * SAIW_RSS_KEEPALIVE_S until the feed follows. Whitespace is allowed before + * the JSON value, and between the XML declaration (sent up front) and <rss>, + * so the body is still a valid document. */ #include <libwebsockets.h> @@ -51,9 +73,24 @@ /* feed readers are asked to wait at least this many minutes between polls */ #define SAIW_RSS_TTL_MINS 5 +/* longest a long poll request is held */ +#define SAIW_RSS_MAX_WAIT_S 600 + +/* interval between keepalive newlines on a held request */ +#define SAIW_RSS_KEEPALIVE_S 10 + +/* + * Held requests at once: the feed is public, and each held request pins a + * connection here and on the front-end proxy. Past this, requests to wait + * are told to retry later. + */ +#define SAIW_RSS_MAX_WAITERS 64 + /* identifies our metadata elements, it needn't resolve to anything */ #define SAIW_RSS_NS "https://warmcat.com/sai/ns/rss" +#define SAIW_RSS_XML_DECL "<?xml version=\"1.0\" encoding=\"UTF-8\"?>\n" + /* * Escaped field sizes: each input byte becomes at most 6 ("&quot;") */ @@ -190,10 +227,10 @@ saiw_rss_col(sqlite3_stmt *sm, int col) } static int -saiw_rss_append(struct pss *pss, const char *buf, int len) +saiw_rss_append(struct pss *pss, const void *buf, size_t len) { - if (len >= 0 && lws_buflist_append_segment(&pss->rss_tx, - (const uint8_t *)buf, (size_t)len) >= 0) + if (lws_buflist_append_segment(&pss->rss_tx, (const uint8_t *)buf, + len) >= 0) return 0; lwsl_err("%s: OOM\n", __func__); @@ -202,51 +239,163 @@ saiw_rss_append(struct pss *pss, const char *buf, int len) } /* + * Collect the events in the feed for this request's scope into f, allocated + * in *ac, and compute f->index from them. This only reads the events table, + * so it's cheap enough to redo for each held request whenever an event + * changes; the task counts, which need each event's own db, are filled in + * separately by saiw_feed_counts() only when the feed is actually sent. + */ +static int +saiw_feed_query(struct vhd *vhd, struct pss *pss, sai_feed_t *f, + struct lwsac **ac) +{ + /* + * The project and branch come from the request url, so they must be + * bound rather than go via lws_struct_sq3_deserialize(), whose filter + * text is spliced into the sql verbatim. + */ + static const char * const q = + "SELECT uuid, repo_name, ref, hash, created, state, " + "ifnull(adhoc,0), repo_fetchurl, weburl FROM events " + "WHERE state != ?3 AND " + "(?1 IS NULL OR repo_name = ?1) AND " + "(?2 IS NULL OR ref = ?2 OR ref = 'refs/heads/' || ?2) AND " + "(?5 IS NULL OR repo_fetchurl = ?5) " + "ORDER BY created DESC LIMIT ?4"; + struct lws_genhash_ctx hc; + uint8_t digest[32]; + sqlite3_stmt *sm = NULL; + sai_feed_item_t *it; + const char *ref; + int rc, ret = 1; + + memset(f, 0, sizeof(*f)); + + if (lws_genhash_init(&hc, LWS_GENHASH_TYPE_SHA256)) + return 1; + + if (sqlite3_prepare_v2(vhd->pdb, q, -1, &sm, NULL) != SQLITE_OK) { + lwsl_err("%s: prepare failed: %s\n", __func__, + sqlite3_errmsg(vhd->pdb)); + goto bail; + } + + /* unbound parameters are NULL, meaning no scoping on that */ + if ((pss->rss_project[0] && sqlite3_bind_text(sm, 1, pss->rss_project, + -1, SQLITE_STATIC)) || + (pss->rss_branch[0] && sqlite3_bind_text(sm, 2, pss->rss_branch, + -1, SQLITE_STATIC)) || + (pss->rss_fetchurl[0] && sqlite3_bind_text(sm, 5, + pss->rss_fetchurl, -1, SQLITE_STATIC)) || + sqlite3_bind_int(sm, 3, SAIES_DELETED) || + sqlite3_bind_int(sm, 4, SAIW_RSS_ITEMS)) { + lwsl_err("%s: bind failed\n", __func__); + goto bail; + } + + while ((rc = sqlite3_step(sm)) == SQLITE_ROW) { + it = lwsac_use_zero(ac, sizeof(*it), 2048); + if (!it) + goto bail; + + lws_strncpy(it->uuid, saiw_rss_col(sm, 0), sizeof(it->uuid)); + lws_strncpy(it->project, saiw_rss_col(sm, 1), + sizeof(it->project)); + lws_strncpy(it->ref, saiw_rss_col(sm, 2), sizeof(it->ref)); + ref = it->ref; + if (!strncmp(ref, "refs/heads/", 11)) + ref += 11; + lws_strncpy(it->branch, ref, sizeof(it->branch)); + lws_strncpy(it->hash, saiw_rss_col(sm, 3), sizeof(it->hash)); + it->received = (uint64_t)sqlite3_column_int64(sm, 4); + it->state = sqlite3_column_int(sm, 5); + it->adhoc = !!sqlite3_column_int(sm, 6); + lws_strncpy(it->fetchurl, saiw_rss_col(sm, 7), + sizeof(it->fetchurl)); + lws_strncpy(it->weburl, saiw_rss_col(sm, 8), + sizeof(it->weburl)); + lws_strncpy(it->state_name, saiw_rss_state_name(it->state), + sizeof(it->state_name)); + + lws_dll2_add_tail(&it->list, &f->items); + + if (lws_genhash_update(&hc, it->uuid, strlen(it->uuid)) || + lws_genhash_update(&hc, ":", 1) || + lws_genhash_update(&hc, it->state_name, + strlen(it->state_name)) || + lws_genhash_update(&hc, "\n", 1)) + goto bail; + } + + if (rc != SQLITE_DONE) { + lwsl_err("%s: step failed: %s\n", __func__, + sqlite3_errmsg(vhd->pdb)); + goto bail; + } + + ret = 0; + +bail: + if (sm) + sqlite3_finalize(sm); + lws_genhash_destroy(&hc, ret ? NULL : digest); + if (!ret) + /* 16 bytes of it is plenty to notice a change */ + lws_hex_from_byte_array(digest, 16, f->index, sizeof(f->index)); + + return ret; +} + +/* the task counts live in each event's own db */ + +static void +saiw_feed_counts(struct vhd *vhd, sai_feed_t *f) +{ + sqlite3 *pdb; + + lws_start_foreach_dll(struct lws_dll2 *, p, f->items.head) { + sai_feed_item_t *it = lws_container_of(p, sai_feed_item_t, + list); + + pdb = NULL; + if (!sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, + vhd->sqlite3_path_lhs, it->uuid, + 0, &pdb)) { + saiw_event_summary_string(pdb, it->uuid, it->summary, + sizeof(it->summary), + &it->tasks_ok, &it->tasks_bad, + &it->tasks_building, + &it->tasks_wait, + &it->tasks_total); + sai_event_db_close(&vhd->sqlite3_cache, &pdb); + } + } lws_end_foreach_dll(p); +} + +/* * Render one event as an <item> and append it to the feed */ static int -saiw_rss_item(struct vhd *vhd, struct pss *pss, sqlite3_stmt *sm) +saiw_rss_item(struct pss *pss, const sai_feed_item_t *it) { - char item[12288], summary[96], esc_sum[SAIW_RSS_ESC(96)], + char item[12288], esc_sum[SAIW_RSS_ESC(96)], esc_uuid[SAIW_RSS_ESC(65)], esc_repo[SAIW_RSS_ESC(65)], esc_branch[SAIW_RSS_ESC(65)], esc_hash[SAIW_RSS_ESC(65)], esc_fetchurl[SAIW_RSS_ESC(96)], esc_weburl[SAIW_RSS_ESC(128)], date[64]; - unsigned int good, bad, ongoing, pending, total; - const char *uuid = saiw_rss_col(sm, 0), *ref = saiw_rss_col(sm, 2), - *state_name; - time_t created = (time_t)sqlite3_column_int64(sm, 4); - int state = sqlite3_column_int(sm, 5), - adhoc = sqlite3_column_int(sm, 6), n; - sqlite3 *pdb = NULL; - - /* the task counts live in the event's own db */ - - summary[0] = '\0'; - good = bad = ongoing = pending = total = 0; - if (!sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, - vhd->sqlite3_path_lhs, uuid, 0, &pdb)) { - saiw_event_summary_string(pdb, uuid, summary, sizeof(summary), - &good, &bad, &ongoing, &pending, - &total); - sai_event_db_close(&vhd->sqlite3_cache, &pdb); - } - - if (!strncmp(ref, "refs/heads/", 11)) - ref += 11; + time_t created = (time_t)it->received; + int n; if (lws_http_date_render_from_unix(date, sizeof(date), &created)) date[0] = '\0'; - state_name = saiw_rss_state_name(state); - - saiw_xml_esc(esc_uuid, sizeof(esc_uuid), uuid); - saiw_xml_esc(esc_repo, sizeof(esc_repo), saiw_rss_col(sm, 1)); - saiw_xml_esc(esc_branch, sizeof(esc_branch), ref); - saiw_xml_esc(esc_hash, sizeof(esc_hash), saiw_rss_col(sm, 3)); - saiw_xml_esc(esc_sum, sizeof(esc_sum), summary); - saiw_xml_esc(esc_fetchurl, sizeof(esc_fetchurl), saiw_rss_col(sm, 7)); - saiw_xml_esc(esc_weburl, sizeof(esc_weburl), saiw_rss_col(sm, 8)); + saiw_xml_esc(esc_uuid, sizeof(esc_uuid), it->uuid); + saiw_xml_esc(esc_repo, sizeof(esc_repo), it->project); + saiw_xml_esc(esc_branch, sizeof(esc_branch), it->branch); + saiw_xml_esc(esc_hash, sizeof(esc_hash), it->hash); + saiw_xml_esc(esc_sum, sizeof(esc_sum), it->summary); + saiw_xml_esc(esc_fetchurl, sizeof(esc_fetchurl), it->fetchurl); + saiw_xml_esc(esc_weburl, sizeof(esc_weburl), it->weburl); n = lws_snprintf(item, sizeof(item), "<item>\n" @@ -271,95 +420,78 @@ saiw_rss_item(struct vhd *vhd, struct pss *pss, sqlite3_stmt *sm) "</item>\n", /* title */ - esc_repo, esc_branch, adhoc ? " (ad-hoc)" : "", state_name, - esc_sum[0] ? " - " : "", esc_sum, + esc_repo, esc_branch, it->adhoc ? " (ad-hoc)" : "", + it->state_name, esc_sum[0] ? " - " : "", esc_sum, /* link */ esc_uuid, /* guid */ - esc_uuid, state_name, + esc_uuid, it->state_name, /* pubDate */ date, /* description */ - adhoc ? "Ad-hoc build of " : "", esc_repo, esc_branch, esc_hash, - state_name, total, total == 1 ? "" : "s", - good, bad, ongoing, pending, + it->adhoc ? "Ad-hoc build of " : "", esc_repo, esc_branch, + esc_hash, it->state_name, it->tasks_total, + it->tasks_total == 1 ? "" : "s", + it->tasks_ok, it->tasks_bad, it->tasks_building, + it->tasks_wait, /* categories */ esc_repo, esc_branch, /* sai: metadata */ - (unsigned long long)created, esc_repo, esc_branch, esc_hash, - esc_fetchurl, esc_weburl, !!adhoc, state, state_name, - total, good, bad, ongoing, pending); + (unsigned long long)it->received, esc_repo, esc_branch, + esc_hash, esc_fetchurl, esc_weburl, it->adhoc, it->state, + it->state_name, it->tasks_total, it->tasks_ok, it->tasks_bad, + it->tasks_building, it->tasks_wait); if (n >= (int)sizeof(item) - 1) { - lwsl_err("%s: item for %s too large\n", __func__, uuid); + lwsl_err("%s: item for %s too large\n", __func__, it->uuid); return 1; } - return saiw_rss_append(pss, item, n); + return saiw_rss_append(pss, item, (size_t)n); } -int -saiw_rss_http(struct vhd *vhd, struct pss *pss, struct lws *wsi) +/* + * Append the feed as RSS to pss->rss_tx. The XML declaration is left out if + * it was already sent at the start of a held request. + */ +static int +saiw_feed_render_xml(struct pss *pss, const sai_feed_t *f, int decl) { - char buf[LWS_PRE + 8192], *start = buf + LWS_PRE, - *end = buf + sizeof(buf) - 1, project[65], branch[65], - uenc_project[SAIW_RSS_ESC(65)], uenc_branch[SAIW_RSS_ESC(65)], - esc_project[SAIW_RSS_ESC(65)], esc_branch[SAIW_RSS_ESC(65)], - self[512], esc_self[SAIW_RSS_ESC(512)], date[64]; - /* - * The query takes the project and branch from the request url, so it - * must bind them rather than go via lws_struct_sq3_deserialize(), - * whose filter text is spliced into the sql verbatim. - */ - static const char * const q = - "SELECT uuid, repo_name, ref, hash, created, state, " - "ifnull(adhoc,0), repo_fetchurl, weburl FROM events " - "WHERE state != ?3 AND " - "(?1 IS NULL OR repo_name = ?1) AND " - "(?2 IS NULL OR ref = ?2 OR ref = 'refs/heads/' || ?2) " - "ORDER BY created DESC LIMIT ?4"; - uint8_t *hs = (uint8_t *)start, *hp = hs, - *he = (uint8_t *)end; - sqlite3_stmt *sm = NULL; + char buf[12288], uenc[SAIW_RSS_ESC(96)], esc_project[SAIW_RSS_ESC(65)], + esc_branch[SAIW_RSS_ESC(65)], self[1024], *sp = self, + *se = self + sizeof(self), esc_self[SAIW_RSS_ESC(1024)], date[64]; + int pl = !!pss->rss_project[0], bl = !!pss->rss_branch[0], n, + first = 1; + const char *scope[][2] = { + { "project", pss->rss_project }, + { "fetchurl", pss->rss_fetchurl }, + { "branch", pss->rss_branch }, + }; time_t now = time(NULL); - int pl, bl, n, rc; - - lws_buflist_destroy_all_segments(&pss->rss_tx); - - pl = lws_get_urlarg_by_name_safe(wsi, "project=", project, - sizeof(project)); - bl = lws_get_urlarg_by_name_safe(wsi, "branch=", branch, - sizeof(branch)); - if (pl == -2 || bl == -2) - /* longer than any project or ref we store */ - return HTTP_STATUS_BAD_REQUEST; - if (pl < 0) - pl = 0; - if (bl < 0) - bl = 0; - project[pl] = '\0'; - branch[bl] = '\0'; + size_t m; /* the feed's own (relative) url, including any scoping */ - lws_urlencode(uenc_project, project, (int)sizeof(uenc_project)); - lws_urlencode(uenc_branch, branch, (int)sizeof(uenc_branch)); - lws_snprintf(self, sizeof(self), "rss.xml%s%s%s%s%s%s", - pl || bl ? "?" : "", - pl ? "project=" : "", pl ? uenc_project : "", - pl && bl ? "&" : "", - bl ? "branch=" : "", bl ? uenc_branch : ""); + sp += lws_snprintf(sp, lws_ptr_diff_size_t(se, sp), "rss.xml"); + for (m = 0; m < LWS_ARRAY_SIZE(scope); m++) { + if (!scope[m][1][0]) + continue; + lws_urlencode(uenc, scope[m][1], (int)sizeof(uenc)); + sp += lws_snprintf(sp, lws_ptr_diff_size_t(se, sp), "%c%s=%s", + first ? '?' : '&', scope[m][0], uenc); + first = 0; + } saiw_xml_esc(esc_self, sizeof(esc_self), self); - saiw_xml_esc(esc_project, sizeof(esc_project), project); - saiw_xml_esc(esc_branch, sizeof(esc_branch), branch); + saiw_xml_esc(esc_project, sizeof(esc_project), pss->rss_project); + saiw_xml_esc(esc_branch, sizeof(esc_branch), pss->rss_branch); if (lws_http_date_render_from_unix(date, sizeof(date), &now)) date[0] = '\0'; - n = lws_snprintf(start, lws_ptr_diff_size_t(end, start), - "<?xml version=\"1.0\" encoding=\"UTF-8\"?>\n" + n = lws_snprintf(buf, sizeof(buf), + "%s" "<rss version=\"2.0\" " "xmlns:atom=\"http://www.w3.org/2005/Atom\" " "xmlns:sai=\"" SAIW_RSS_NS "\">\n" @@ -373,83 +505,306 @@ saiw_rss_http(struct vhd *vhd, struct pss *pss, struct lws *wsi) "type=\"application/rss+xml\"/>\n" "<generator>Sai</generator>\n" "<lastBuildDate>%s</lastBuildDate>\n" - "<ttl>%d</ttl>\n", + "<ttl>%d</ttl>\n" + "<sai:index>%s</sai:index>\n", + decl ? SAIW_RSS_XML_DECL : "", pl ? ": " : "", esc_project, bl ? (pl ? " " : ": ") : "", esc_branch, SAIW_RSS_ITEMS, pl ? " for " : "", esc_project, bl ? " branch " : "", esc_branch, - esc_self, date, SAIW_RSS_TTL_MINS); + esc_self, date, SAIW_RSS_TTL_MINS, f->index); - if (n >= lws_ptr_diff(end, start) - 1) { + if (n >= (int)sizeof(buf) - 1) { lwsl_err("%s: channel header too large\n", __func__); - goto bail; + + return 1; } - if (saiw_rss_append(pss, start, n)) - goto bail; + if (saiw_rss_append(pss, buf, (size_t)n)) + return 1; - if (sqlite3_prepare_v2(vhd->pdb, q, -1, &sm, NULL) != SQLITE_OK) { - lwsl_err("%s: prepare failed: %s\n", __func__, - sqlite3_errmsg(vhd->pdb)); - goto bail; - } + lws_start_foreach_dll(struct lws_dll2 *, p, f->items.head) { + if (saiw_rss_item(pss, lws_container_of(p, sai_feed_item_t, + list))) + return 1; + } lws_end_foreach_dll(p); - /* unbound parameters are NULL, meaning no scoping on that */ - if ((pl && sqlite3_bind_text(sm, 1, project, pl, SQLITE_STATIC)) || - (bl && sqlite3_bind_text(sm, 2, branch, bl, SQLITE_STATIC)) || - sqlite3_bind_int(sm, 3, SAIES_DELETED) || - sqlite3_bind_int(sm, 4, SAIW_RSS_ITEMS)) { - lwsl_err("%s: bind failed\n", __func__); - goto bail; + return saiw_rss_append(pss, "</channel>\n</rss>\n", 18); +} + +static int +saiw_feed_render_json(struct pss *pss, sai_feed_t *f) +{ + lws_struct_serialize_t *js; + uint8_t buf[4096]; + size_t w; + int r; + + js = lws_struct_json_serialize_create(lsm_schema_json_map_feed, + LWS_ARRAY_SIZE(lsm_schema_json_map_feed), 0, f); + if (!js) + return 1; + + do { + w = 0; + r = lws_struct_json_serialize(js, buf, sizeof(buf), &w); + if (r == LSJS_RESULT_ERROR || + (w && saiw_rss_append(pss, buf, w))) { + lws_struct_json_serialize_destroy(&js); + + return 1; + } + } while (r == LSJS_RESULT_CONTINUE); + + lws_struct_json_serialize_destroy(&js); + + return saiw_rss_append(pss, "\n", 1); +} + +/* + * Fill in the counts and append the whole feed body for this request + */ +static int +saiw_feed_render(struct vhd *vhd, struct pss *pss, sai_feed_t *f, int decl) +{ + saiw_feed_counts(vhd, f); + + if (pss->rss_json) + return saiw_feed_render_json(pss, f); + + return saiw_feed_render_xml(pss, f, decl); +} + +static void +saiw_rss_unpark(struct pss *pss) +{ + lws_dll2_remove(&pss->rss_list); + lws_sul_cancel(&pss->sul_rss); +} + +/* + * A held request is being answered: send the current feed to finish it + */ +static void +saiw_rss_finish(struct vhd *vhd, struct pss *pss) +{ + struct lwsac *ac = NULL; + sai_feed_t f; + + saiw_rss_unpark(pss); + + /* + * The status and headers went out long ago, so a failure now can only + * be reported by dropping the connection, which the client handles + * like any other broken connection + */ + pss->rss_state = SAIW_RSS_FAILED; + if (!saiw_feed_query(vhd, pss, &f, &ac) && + !saiw_feed_render(vhd, pss, &f, 0)) + pss->rss_state = SAIW_RSS_FINISHING; + + lwsac_free(&ac); + lws_callback_on_writable(pss->wsi); +} + +static void +saiw_rss_sul_cb(lws_sorted_usec_list_t *sul) +{ + struct pss *pss = lws_container_of(sul, struct pss, sul_rss); + lws_usec_t now = lws_now_usecs(), next; + + if (now >= pss->rss_deadline) { + /* waited as long as asked, answer with the current feed */ + saiw_rss_finish(pss->vhd, pss); + + return; } - while ((rc = sqlite3_step(sm)) == SQLITE_ROW) - if (saiw_rss_item(vhd, pss, sm)) - goto bail; + pss->rss_ka = 1; + lws_callback_on_writable(pss->wsi); - if (rc != SQLITE_DONE) { - lwsl_err("%s: step failed: %s\n", __func__, - sqlite3_errmsg(vhd->pdb)); + next = SAIW_RSS_KEEPALIVE_S * LWS_US_PER_SEC; + if (pss->rss_deadline - now < next) + next = pss->rss_deadline - now; + + lws_sul_schedule(pss->vhd->context, 0, &pss->sul_rss, saiw_rss_sul_cb, + next); +} + +/* + * Some event joined, left or changed state: answer any held requests whose + * feed index is now out of date + */ +void +saiw_rss_event_change(struct vhd *vhd) +{ + lws_start_foreach_dll_safe(struct lws_dll2 *, p, p1, + vhd->rss_waiters.head) { + struct pss *pss = lws_container_of(p, struct pss, rss_list); + struct lwsac *ac = NULL; + sai_feed_t f; + + if (!saiw_feed_query(vhd, pss, &f, &ac) && + strcmp(f.index, pss->rss_index)) + saiw_rss_finish(vhd, pss); + + lwsac_free(&ac); + } lws_end_foreach_dll_safe(p, p1); +} + +/* + * Copy a url arg into buf, "" if absent. Returns nonzero if it was present + * but too long for buf. + */ +static int +saiw_rss_urlarg(struct lws *wsi, const char *name, char *buf, size_t len) +{ + int n = lws_get_urlarg_by_name_safe(wsi, name, buf, (int)len); + + if (n < 0) + buf[0] = '\0'; + + return n == -2; +} + +/* + * Returns < 0 to close the connection, 0 if the response is under way, 1 if + * the response was completed here, or else an http status for the caller to + * respond with + */ +int +saiw_rss_http(struct vhd *vhd, struct pss *pss, struct lws *wsi, int json) +{ + char buf[LWS_PRE + 1024], index[33], wait_s[12]; + uint8_t *start = (uint8_t *)buf + LWS_PRE, *p = start, + *end = (uint8_t *)buf + sizeof(buf) - 1; + const char *ctype = json ? "application/json" : + "application/rss+xml; charset=utf-8"; + struct lwsac *ac = NULL; + sai_feed_t f; + int wait; + + saiw_rss_close(pss); + pss->wsi = wsi; + pss->rss_json = (uint8_t)!!json; + + if (saiw_rss_urlarg(wsi, "project=", pss->rss_project, + sizeof(pss->rss_project)) || + saiw_rss_urlarg(wsi, "branch=", pss->rss_branch, + sizeof(pss->rss_branch)) || + saiw_rss_urlarg(wsi, "fetchurl=", pss->rss_fetchurl, + sizeof(pss->rss_fetchurl))) + /* longer than any project or ref we store */ + return HTTP_STATUS_BAD_REQUEST; + + if (saiw_rss_urlarg(wsi, "index=", index, sizeof(index))) + /* can't be a token of ours, so it's out of date */ + index[0] = '\0'; + + wait = 0; + if (!saiw_rss_urlarg(wsi, "wait=", wait_s, sizeof(wait_s))) + wait = atoi(wait_s); + if (wait < 0) + wait = 0; + if (wait > SAIW_RSS_MAX_WAIT_S) + wait = SAIW_RSS_MAX_WAIT_S; + + if (saiw_feed_query(vhd, pss, &f, &ac)) goto bail; + + if (wait && index[0] && !strcmp(index, f.index)) { + + /* nothing new for them yet: hold the request */ + + lwsac_free(&ac); + + if (vhd->rss_waiters.count >= SAIW_RSS_MAX_WAITERS) { + lwsl_notice("%s: %u waiting, deferring another\n", + __func__, + (unsigned int)vhd->rss_waiters.count); + + if (lws_add_http_header_status(wsi, + HTTP_STATUS_SERVICE_UNAVAILABLE, + &p, end) || + lws_add_http_header_by_token(wsi, + WSI_TOKEN_HTTP_RETRY_AFTER, + (const uint8_t *)"60", 2, &p, end) || + lws_add_http_header_content_length(wsi, 0, &p, end) || + lws_finalize_write_http_header(wsi, start, &p, end)) + return -1; + + return 1; + } + + if (lws_add_http_common_headers(wsi, HTTP_STATUS_OK, ctype, + LWS_ILLEGAL_HTTP_CONTENT_LEN, + &p, end) || + lws_add_http_header_by_token(wsi, + WSI_TOKEN_HTTP_CACHE_CONTROL, + (const uint8_t *)"no-store", 8, &p, end) || + lws_finalize_write_http_header(wsi, start, &p, end)) + return -1; + + /* + * We are a long-lived response now: this lifts the http + * timeouts (and for h2, the network connection's idle one) + */ + lws_http_mark_sse(wsi); + + if (!json && saiw_rss_append(pss, SAIW_RSS_XML_DECL, + strlen(SAIW_RSS_XML_DECL))) + return -1; + + lws_strncpy(pss->rss_index, index, sizeof(pss->rss_index)); + pss->rss_state = SAIW_RSS_PARKED; + pss->rss_deadline = lws_now_usecs() + + (lws_usec_t)wait * LWS_US_PER_SEC; + lws_dll2_add_tail(&pss->rss_list, &vhd->rss_waiters); + lws_sul_schedule(vhd->context, 0, &pss->sul_rss, + saiw_rss_sul_cb, + (wait < SAIW_RSS_KEEPALIVE_S ? wait : + SAIW_RSS_KEEPALIVE_S) * LWS_US_PER_SEC); + + lws_callback_on_writable(wsi); + + return 0; } - sqlite3_finalize(sm); - sm = NULL; + /* answer straight away */ - n = lws_snprintf(start, lws_ptr_diff_size_t(end, start), - "</channel>\n</rss>\n"); - if (saiw_rss_append(pss, start, n)) + if (saiw_feed_render(vhd, pss, &f, 1)) goto bail; + lwsac_free(&ac); + /* * The whole feed is rendered, so we know the length. The results * change as builds progress, so caches may only hold it briefly. */ - if (lws_add_http_header_status(wsi, HTTP_STATUS_OK, &hp, he) || + if (lws_add_http_header_status(wsi, HTTP_STATUS_OK, &p, end) || lws_add_http_header_content_length(wsi, (lws_filepos_t)lws_buflist_total_len(&pss->rss_tx), - &hp, he) || + &p, end) || lws_add_http_header_by_token(wsi, WSI_TOKEN_HTTP_CONTENT_TYPE, - (const uint8_t *)"application/rss+xml; charset=utf-8", - 34, &hp, he) || + (const uint8_t *)ctype, (int)strlen(ctype), &p, end) || lws_add_http_header_by_token(wsi, WSI_TOKEN_HTTP_CACHE_CONTROL, - (const uint8_t *)"public, max-age=60", 18, &hp, he) || - lws_finalize_write_http_header(wsi, hs, &hp, he)) { - lws_buflist_destroy_all_segments(&pss->rss_tx); + (const uint8_t *)"public, max-age=60", 18, &p, end) || + lws_finalize_write_http_header(wsi, start, &p, end)) { + saiw_rss_close(pss); return -1; } + pss->rss_state = SAIW_RSS_FINISHING; lws_callback_on_writable(wsi); return 0; bail: - if (sm) - sqlite3_finalize(sm); - lws_buflist_destroy_all_segments(&pss->rss_tx); + lwsac_free(&ac); + saiw_rss_close(pss); return HTTP_STATUS_INTERNAL_SERVER_ERROR; } @@ -458,25 +813,46 @@ int saiw_rss_writeable(struct pss *pss, struct lws *wsi) { uint8_t buf[LWS_PRE + 4096], *seg; - size_t n = lws_buflist_next_segment_len(&pss->rss_tx, &seg); + size_t n; int final; + if (pss->rss_state == SAIW_RSS_FAILED) + return -1; + + n = lws_buflist_next_segment_len(&pss->rss_tx, &seg); + if (!n) { + /* a held request with nothing to send but a keepalive */ + if (pss->rss_state != SAIW_RSS_PARKED || !pss->rss_ka) + return 0; + + pss->rss_ka = 0; + buf[LWS_PRE] = '\n'; + + return lws_write(wsi, buf + LWS_PRE, 1, LWS_WRITE_HTTP) != 1 ? + -1 : 0; + } + if (n > sizeof(buf) - LWS_PRE) n = sizeof(buf) - LWS_PRE; memcpy(buf + LWS_PRE, seg, n); lws_buflist_use_segment(&pss->rss_tx, n); - final = !lws_buflist_total_len(&pss->rss_tx); + final = pss->rss_state == SAIW_RSS_FINISHING && + !lws_buflist_total_len(&pss->rss_tx); if (lws_write(wsi, buf + LWS_PRE, n, final ? LWS_WRITE_HTTP_FINAL : LWS_WRITE_HTTP) != (int)n) return -1; if (!final) { - lws_callback_on_writable(wsi); + if (lws_buflist_total_len(&pss->rss_tx)) + lws_callback_on_writable(wsi); return 0; } + pss->rss_state = SAIW_RSS_IDLE; + + /* a held request's response was close-delimited on h1 */ if (lws_http_transaction_completed(wsi)) return -1; @@ -487,6 +863,11 @@ void saiw_rss_close(struct pss *pss) { /* NULL if the conn went before per-session storage was allocated */ - if (pss) - lws_buflist_destroy_all_segments(&pss->rss_tx); + if (!pss) + return; + + saiw_rss_unpark(pss); + lws_buflist_destroy_all_segments(&pss->rss_tx); + pss->rss_state = SAIW_RSS_IDLE; + pss->rss_ka = 0; } diff --git a/src/web/w-ws-browser.c b/src/web/w-ws-browser.c index 94684b8..b4b2f03 100644 --- a/src/web/w-ws-browser.c +++ b/src/web/w-ws-browser.c @@ -850,6 +850,9 @@ saiw_browsers_task_state_change(struct vhd *vhd, const char *task_uuid) int saiw_event_state_change(struct vhd *vhd, const char *event_uuid) { + /* long poll feed requests may be waiting on this */ + saiw_rss_event_change(vhd); + lws_start_foreach_dll(struct lws_dll2 *, p, vhd->browsers.head) { struct pss *pss = lws_container_of(p, struct pss, same);
Page fetched 0s ago, creation time: 29ms (vhost etag hits: 0%, cache hits: 0%)