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-aarch64-apple-m1.svg
Author[]Andy Green <andy@warmcat.com> 2026-05-17 08:58 UTC
Committer[]Andy Green <andy@warmcat.com> 2026-05-17 09:02 UTC
Tree3f314596fe7553dbf7a652f6631f1f53701a98ad   Raw Patch
 
sai-virt: phase1
sai-virt: phase1
diff --git a/src/common/include/private.h b/src/common/include/private.h index 44e7bdc..b4773c8 100644 --- a/src/common/include/private.h +++ b/src/common/include/private.h @@ -755,6 +755,18 @@ typedef struct sai_pcon_control { char on; } sai_pcon_control_t; +typedef struct sai_platform_pending_task { + lws_dll2_t list; + char plat[64]; + unsigned int pending; +} sai_platform_pending_task_t; + +typedef struct sai_platform_pending_tasks { + lws_dll2_t list; + char pcons[1024]; + lws_dll2_owner_t tasks; /* sai_platform_pending_task_t */ +} sai_platform_pending_tasks_t; + /* * Because the definitions of these arrays of map structs are mostly in * common/struct-metadata.c, we are forced to repeat the length of the struct @@ -821,7 +833,10 @@ extern const lws_struct_map_t lsm_watcher[9], lsm_schema_sq3_map_watcher[1], lsm_schema_json_map_watcher[1], - lsm_watcher_conf[1]; + lsm_watcher_conf[1], + lsm_pending_task[2], + lsm_pending_tasks[2], + lsm_schema_pending_tasks[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 34e4cd9..a07826d 100644 --- a/src/common/struct-metadata.c +++ b/src/common/struct-metadata.c @@ -516,3 +516,20 @@ 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"), }; + +const lws_struct_map_t lsm_pending_task[] = { + LSM_CARRAY (sai_platform_pending_task_t, plat, "plat"), + LSM_UNSIGNED (sai_platform_pending_task_t, pending, "pending"), +}; + +const lws_struct_map_t lsm_pending_tasks[] = { + LSM_CARRAY (sai_platform_pending_tasks_t, pcons, "pcons"), + LSM_LIST (sai_platform_pending_tasks_t, tasks, + sai_platform_pending_task_t, list, + NULL, lsm_pending_task, "tasks"), +}; + +const lws_struct_map_t lsm_schema_pending_tasks[] = { + LSM_SCHEMA(sai_platform_pending_tasks_t, NULL, lsm_pending_tasks, + "com.warmcat.sai.power.pending_tasks"), +}; diff --git a/src/power/p-ws-server.c b/src/power/p-ws-server.c index f745071..02403b1 100644 --- a/src/power/p-ws-server.c +++ b/src/power/p-ws-server.c @@ -36,12 +36,13 @@ static const lws_struct_map_t lsm_schema_power_state[] = { "com.warmcat.sai.powerstate"), }; -/* Combined map for RX from server */ static const lws_struct_map_t lsm_saip_rx_map[] = { LSM_SCHEMA(sai_stay_t, NULL, lsm_stay, "com.warmcat.sai.power.stay"), LSM_SCHEMA(sai_pcon_control_t, NULL, lsm_pcon_control, "com.warmcat.sai.pcon_control"), + LSM_SCHEMA(sai_platform_pending_tasks_t, NULL, lsm_pending_tasks, + "com.warmcat.sai.power.pending_tasks"), }; int @@ -236,7 +237,7 @@ saip_m_rx(void *userobj, const uint8_t *buf, size_t len, int flags) } else { lwsl_warn("%s: Unknown PCON '%s'\n", __func__, ctl->pcon_name); } - } else { + } else if (a.top_schema_index == 0) { /* Stay */ sai_stay_t *stay = (sai_stay_t *)a.dest; saip_pcon_t *pc; @@ -299,94 +300,94 @@ saip_m_rx(void *userobj, const uint8_t *buf, size_t len, int flags) } } lws_end_foreach_dll(b_node); } lws_end_foreach_dll(p); - } + } else if (a.top_schema_index == 2) { + /* Pending tasks JSON */ + sai_platform_pending_tasks_t *pt = (sai_platform_pending_tasks_t *)a.dest; -found: - lwsac_free(&a.ac); - return 0; - } - lwsac_free(&a.ac); + lwsl_info("%s: received pending tasks pcons: %s\n", __func__, pt->pcons); - /* - * It wasn't a JSON message... it's the comma-separated list of needed - * platforms then - */ + lws_start_foreach_dll(struct lws_dll2 *, p, power.sai_pcon_owner.head) { + saip_pcon_t *pc = lws_container_of(p, saip_pcon_t, list); - lwsl_err("%s: ************* Received comma-separated list of needed platforms\n", __func__); - sai_dump_stderr(buf, len); + pc->flags &= (uint8_t)~SAIP_PCON_F_NEEDED; + } lws_end_foreach_dll(p); - lws_start_foreach_dll(struct lws_dll2 *, p, power.sai_pcon_owner.head) { - saip_pcon_t *pc = lws_container_of(p, saip_pcon_t, list); + if (pt->pcons[0]) { + const char *cp = pt->pcons; + const char *end = cp + strlen(cp); + + while (cp < end) { + const char *comma = memchr(cp, ',', lws_ptr_diff_size_t(end, cp)); + size_t token_len; + + if (comma) + token_len = lws_ptr_diff_size_t(comma, cp); + else + token_len = lws_ptr_diff_size_t(end, cp); + + if (token_len) { + char pcon[64]; + saip_pcon_t *pc; + + lws_strnncpy(pcon, cp, token_len, sizeof(pcon)); + pc = saip_pcon_by_name(&power, pcon); + if (pc) + pc->flags |= SAIP_PCON_F_NEEDED; + else + lwsl_notice("%s: unknown pcon '%.*s' needed\n", + __func__, (int)token_len, cp); + } - pc->flags &= (uint8_t)~SAIP_PCON_F_NEEDED; - } lws_end_foreach_dll(p); + cp += token_len; + if (cp < end && *cp == ',') + cp++; + } + } - if (len) { - const char *cp = (const char *)buf; - const char *end = cp + len; - - while (cp < end) { - const char *comma = memchr(cp, ',', lws_ptr_diff_size_t(end, cp)); - size_t token_len; - - if (comma) - token_len = lws_ptr_diff_size_t(comma, cp); - else - token_len = lws_ptr_diff_size_t(end, cp); - - if (token_len) { - char pcon[64]; - saip_pcon_t *pc; - - lws_strnncpy(pcon, cp, token_len, sizeof(pcon)); - pc = saip_pcon_by_name(&power, pcon); - if (pc) - pc->flags |= SAIP_PCON_F_NEEDED; - else - lwsl_notice("%s: unknown pcon '%.*s' needed\n", - __func__, (int)token_len, cp); + /* + * Propagate needed state up the dependency tree + * + * If a PCON is needed, and it depends on another PCON, that parent PCON + * is also needed. + */ + { + int changed; + + do { + changed = 0; + lws_start_foreach_dll(struct lws_dll2 *, p, + power.sai_pcon_owner.head) { + saip_pcon_t *pc = lws_container_of(p, + saip_pcon_t, list); + + if (pc->flags & SAIP_PCON_F_NEEDED) { + /* check if this PCON depends on another */ + if (pc->depends_on) { + saip_pcon_t *parent = saip_pcon_by_name(&power, + pc->depends_on); + if (parent && !(parent->flags & SAIP_PCON_F_NEEDED)) { + parent->flags |= SAIP_PCON_F_NEEDED; + changed = 1; + lwsl_notice("%s: PCON %s needed by dep %s\n", + __func__, parent->name, pc->name); + } + } + } + } lws_end_foreach_dll(p); + } while (changed); } - cp += token_len; - if (cp < end && *cp == ',') - cp++; + saip_pcon_start_check(); } - } - /* - * Propagate needed state up the dependency tree - * - * If a PCON is needed, and it depends on another PCON, that parent PCON - * is also needed. - */ - { - int changed; - - do { - changed = 0; - lws_start_foreach_dll(struct lws_dll2 *, p, - power.sai_pcon_owner.head) { - saip_pcon_t *pc = lws_container_of(p, - saip_pcon_t, list); - - if (pc->flags & SAIP_PCON_F_NEEDED) { - /* check if this PCON depends on another */ - if (pc->depends_on) { - saip_pcon_t *parent = saip_pcon_by_name(&power, - pc->depends_on); - if (parent && !(parent->flags & SAIP_PCON_F_NEEDED)) { - parent->flags |= SAIP_PCON_F_NEEDED; - changed = 1; - lwsl_notice("%s: PCON %s needed by dep %s\n", - __func__, parent->name, pc->name); - } - } - } - } lws_end_foreach_dll(p); - } while (changed); +found: + lwsac_free(&a.ac); + return 0; } + lwsac_free(&a.ac); - saip_pcon_start_check(); + lwsl_err("%s: ************* Received unknown non-JSON message\n", __func__); + sai_dump_stderr(buf, len); return 0; } diff --git a/src/server/s-power.c b/src/server/s-power.c index ef0b643..1d77d5a 100644 --- a/src/server/s-power.c +++ b/src/server/s-power.c @@ -423,7 +423,46 @@ sais_power_tx(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t bl) } lws_end_foreach_dll(px); - lwsl_info("%s: final pcon list: '%.*s'\n", __func__, + { + sai_platform_pending_tasks_t pt; + lws_struct_serialize_t *js; + struct lwsac *ac = NULL; + + memset(&pt, 0, sizeof(pt)); + + /* copy the pcon list into the struct */ + if (lws_ptr_diff_size_t(p, start) < sizeof(pt.pcons)) + lws_strncpy(pt.pcons, (const char *)start, lws_ptr_diff_size_t(p, start) + 1); + else + lws_strncpy(pt.pcons, (const char *)start, sizeof(pt.pcons)); + + lws_start_foreach_dll(struct lws_dll2 *, px1, vhd->pending_plats.head) { + sais_plat_t *pl = lws_container_of(px1, sais_plat_t, list); + sai_platform_pending_task_t *ptask = lwsac_use_zero(&ac, sizeof(*ptask), 1024); + + if (ptask) { + lws_strncpy(ptask->plat, pl->plat, sizeof(ptask->plat)); + ptask->pending = (unsigned int)pl->pending_count; + lws_dll2_add_tail(&ptask->list, &pt.tasks); + } + } lws_end_foreach_dll(px1); + + js = lws_struct_json_serialize_create(lsm_schema_pending_tasks, + LWS_ARRAY_SIZE(lsm_schema_pending_tasks), 0, &pt); + if (!js) { + lwsl_err("%s: failed to serialize pending tasks\n", __func__); + lwsac_free(&ac); + return -1; + } + + n = (int)lws_struct_json_serialize(js, start, lws_ptr_diff_size_t(end, start), &w); + lws_struct_json_serialize_destroy(&js); + lwsac_free(&ac); + + p = start + w; + } + + lwsl_info("%s: final json: '%.*s'\n", __func__, (int)lws_ptr_diff_size_t(p, start), start); /* diff --git a/src/server/s-private.h b/src/server/s-private.h index 9f6c62d..9e47cef 100644 --- a/src/server/s-private.h +++ b/src/server/s-private.h @@ -205,6 +205,7 @@ typedef struct sais_plat { lws_dll2_t list; const char *plat; char busy; + int pending_count; } sais_plat_t; struct vhd { diff --git a/src/server/s-task.c b/src/server/s-task.c index 704b561..c2a2668 100644 --- a/src/server/s-task.c +++ b/src/server/s-task.c @@ -444,15 +444,17 @@ bail: */ static int -sais_find_or_add_pending_plat(struct vhd *vhd, const char *name) +sais_find_or_add_pending_plat(struct vhd *vhd, const char *name, int count) { sais_plat_t *sp; lws_start_foreach_dll(struct lws_dll2 *, p, vhd->pending_plats.head) { sais_plat_t *pl = lws_container_of(p, sais_plat_t, list); - if (!strcmp(pl->plat, name)) + if (!strcmp(pl->plat, name)) { + pl->pending_count += count; return 1; + } } lws_end_foreach_dll(p); @@ -462,6 +464,7 @@ sais_find_or_add_pending_plat(struct vhd *vhd, const char *name) sp->plat = (const char *)&sp[1]; /* start of overcommit */ memcpy(&sp[1], name, strlen(name) + 1); + sp->pending_count = count; lws_dll2_add_tail(&sp->list, &vhd->pending_plats); @@ -533,9 +536,9 @@ sais_platforms_with_tasks_pending(struct vhd *vhd) if (!sai_event_db_ensure_open(vhd->context, &vhd->sqlite3_cache, vhd->sqlite3_path_lhs, e->uuid, 0, &pdb)) { - if (sqlite3_prepare_v2(pdb, "select distinct platform " + if (sqlite3_prepare_v2(pdb, "select platform, count(*) " "from tasks where " - "(state = 0 or state = 1 or state = 2)", -1, &sm, + "(state = 0 or state = 1 or state = 2) group by platform", -1, &sm, NULL) != SQLITE_OK) { lwsl_err("%s: Unable to %s\n", __func__, sqlite3_errmsg(pdb)); @@ -547,7 +550,8 @@ sais_platforms_with_tasks_pending(struct vhd *vhd) n = sqlite3_step(sm); if (n == SQLITE_ROW) sais_find_or_add_pending_plat(vhd, - (const char *)sqlite3_column_text(sm, 0)); + (const char *)sqlite3_column_text(sm, 0), + sqlite3_column_int(sm, 1)); } while (n == SQLITE_ROW); sqlite3_reset(sm);
Page fetched 0s ago, creation time: 5ms (vhost etag hits: 0%, cache hits: 0%)