| Author | Andy Green <andy@warmcat.com> 2026-09-28 15:27 UTC | | Committer | Andy Green <andy@warmcat.com> 2026-09-28 15:29 UTC | | Tree | 4fdc1f671e44d096fa8bd0bc4826e27c2eb22f89 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

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 (""")
*/
@@ -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);
|