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 / failed-login.html
Author[]Andy Green <andy@warmcat.com> 2025-12-27 08:24 UTC
Committer[]Andy Green <andy@warmcat.com> 2025-12-27 08:43 UTC
Tree92cfe34ae8c1366e12391841180aa66fd70bbd6e   Raw Patch
 
Implement PCON energy monitoring and control.
Implement PCON energy monitoring and control.

- **Backend (sai-power):** Added logic to poll Tasmota PCONs for energy stats (`?m=1`), parse the response, and queue `com.warmcat.sai.pcon_energy` reports.
- **Backend (sai-server):** Updated to forward `com.warmcat.sai.pcon_energy` messages to `sai-web` and route `com.warmcat.sai.pcon_control` messages from `sai-web` to `sai-power`.
- **Frontend (sai-web):** Added handlers to display real-time energy stats in the PCON header and a context menu to toggle PCON power.
- **Security:** Used `textContent` for rendering energy stats to prevent XSS.

Co-developed-by: Gemini 3.0 Pro
diff --git a/assets/sai.js b/assets/sai.js index 9f545f0..7d976c9 100644 --- a/assets/sai.js +++ b/assets/sai.js @@ -1413,6 +1413,34 @@ function createPconDiv(pcon) { { label: `<b>PCON:</b> ${pcon.name}` } ]; + if (authd) { + if (pcon.on) { + menuItems.push({ + label: "Turn Off", + callback: () => { + const msg = { + schema: "com.warmcat.sai.pcon_control", + pcon_name: pcon.name, + on: 0 + }; + sai.send(JSON.stringify(msg)); + } + }); + } else { + menuItems.push({ + label: "Turn On", + callback: () => { + const msg = { + schema: "com.warmcat.sai.pcon_control", + pcon_name: pcon.name, + on: 1 + }; + sai.send(JSON.stringify(msg)); + } + }); + } + } + header.addEventListener("contextmenu", function(event) { if (!authd) return; createContextMenu(event, menuItems); @@ -1702,6 +1730,35 @@ function ws_open_sai() } break; + case "com.warmcat.sai.pcon_energy": + if (jso.items) { + jso.items.forEach(item => { + const pconDiv = document.getElementById("pcon-" + item.name); + if (pconDiv) { + let header = pconDiv.querySelector(".pcon-header"); + let stats = header.querySelector(".pcon-stats"); + if (!stats) { + stats = document.createElement("span"); + stats.className = "pcon-stats"; + stats.style.marginLeft = "10px"; + stats.style.fontSize = "0.9em"; + stats.style.color = "#666"; + header.appendChild(stats); + } + + const d = item; +// stats.textContent = `${d.voltage_v}V ${d.active_power_w}W ${d.current_ma}mA today:${(d.energy_today_wh/1000).toFixed(3)}kWh`; + if (d.voltage_v < 70) + stats.textContent = "unpowered"; + else if (!d.active_power_w) + stats.textContent = "OFF"; + else + stats.textContent = `${d.active_power_w}W`; + } + }); + } + break; + case "com.warmcat.sai.build-metric": var summaryDiv = document.getElementById("metrics-summary-" + jso.task_uuid); if (summaryDiv) { diff --git a/src/common/include/private.h b/src/common/include/private.h index 9adf49b..d887338 100644 --- a/src/common/include/private.h +++ b/src/common/include/private.h @@ -649,6 +649,33 @@ typedef struct sai_builder_registration { char power_controller_name[64]; } sai_builder_registration_t; +typedef struct tasmota_data { + unsigned int voltage_v; + unsigned int current_ma; + unsigned int active_power_w; + unsigned int apparent_power_va; + unsigned int reactive_power_var; + unsigned int power_factor_scaled_1000; + unsigned int energy_today_wh; + unsigned int energy_yesterday_wh; + unsigned int energy_total_wh; +} tasmota_data_t; + +typedef struct sai_pcon_energy_report_item { + lws_dll2_t list; + tasmota_data_t data; + char name[64]; +} sai_pcon_energy_report_item_t; + +typedef struct sai_pcon_energy_report { + lws_dll2_owner_t items; +} sai_pcon_energy_report_t; + +typedef struct sai_pcon_control { + lws_dll2_t list; + char pcon_name[64]; + char on; +} sai_pcon_control_t; /* * Because the definitions of these arrays of map structs are mostly in @@ -704,7 +731,11 @@ extern const lws_struct_map_t lsm_builder_registration[3], lsm_schema_sq3_map_power_controller[1], lsm_schema_sq3_map_controlled_builder[1], - lsm_schema_builder_registration[1]; + lsm_schema_builder_registration[1], + lsm_pcon_energy_report[1], + lsm_schema_pcon_energy[1], + lsm_pcon_control[2], + lsm_schema_pcon_control[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 516d39a..8f3eed1 100644 --- a/src/common/struct-metadata.c +++ b/src/common/struct-metadata.c @@ -393,3 +393,37 @@ const lws_struct_map_t lsm_schema_builder_registration[] = { lsm_builder_registration, "com.warmcat.sai.builder_registration"), }; + +static const lws_struct_map_t lsm_pcon_energy_item[] = { + LSM_CARRAY (sai_pcon_energy_report_item_t, name, "name"), + LSM_UNSIGNED (sai_pcon_energy_report_item_t, data.voltage_v, "voltage_v"), + LSM_UNSIGNED (sai_pcon_energy_report_item_t, data.current_ma, "current_ma"), + LSM_UNSIGNED (sai_pcon_energy_report_item_t, data.active_power_w, "active_power_w"), + LSM_UNSIGNED (sai_pcon_energy_report_item_t, data.apparent_power_va, "apparent_power_va"), + LSM_UNSIGNED (sai_pcon_energy_report_item_t, data.reactive_power_var,"reactive_power_var"), + LSM_UNSIGNED (sai_pcon_energy_report_item_t, data.power_factor_scaled_1000, "power_factor_scaled_1000"), + LSM_UNSIGNED (sai_pcon_energy_report_item_t, data.energy_today_wh, "energy_today_wh"), + LSM_UNSIGNED (sai_pcon_energy_report_item_t, data.energy_yesterday_wh,"energy_yesterday_wh"), + LSM_UNSIGNED (sai_pcon_energy_report_item_t, data.energy_total_wh, "energy_total_wh"), +}; + +const lws_struct_map_t lsm_pcon_energy_report[] = { + LSM_LIST (sai_pcon_energy_report_t, items, + sai_pcon_energy_report_item_t, list, + NULL, lsm_pcon_energy_item, "items"), +}; + +const lws_struct_map_t lsm_schema_pcon_energy[] = { + LSM_SCHEMA(sai_pcon_energy_report_t, NULL, lsm_pcon_energy_report, + "com.warmcat.sai.pcon_energy"), +}; + +const lws_struct_map_t lsm_pcon_control[] = { + LSM_CARRAY (sai_pcon_control_t, pcon_name, "pcon_name"), + LSM_UNSIGNED (sai_pcon_control_t, on, "on"), +}; + +const lws_struct_map_t lsm_schema_pcon_control[] = { + LSM_SCHEMA(sai_pcon_control_t, NULL, lsm_pcon_control, + "com.warmcat.sai.pcon_control"), +}; diff --git a/src/power/p-private.h b/src/power/p-private.h index 345c96d..c401f88 100644 --- a/src/power/p-private.h +++ b/src/power/p-private.h @@ -43,18 +43,6 @@ #define SAI_POWERDOWN_HOLDOFF_US (50 * LWS_US_PER_SEC) -typedef struct tasmota_data { - unsigned int voltage_v; - unsigned int current_ma; - unsigned int active_power_w; - unsigned int apparent_power_va; - unsigned int reactive_power_var; - unsigned int power_factor_scaled_1000; - unsigned int energy_today_wh; - unsigned int energy_yesterday_wh; - unsigned int energy_total_wh; -} tasmota_data_t; - typedef struct tasmota_parse { tasmota_data_t td; struct lws_tokenize ts; @@ -88,6 +76,13 @@ typedef struct saip_pcon { struct lws_ss_handle *ss_tasmota_off; struct lws_ss_handle *ss_tasmota_monitor; + tasmota_data_t latest_data; + lws_usec_t last_monitor_time; + + /* For RX accumulation */ + char monitor_rx_buf[4096]; + size_t monitor_rx_pos; + char on; char manual_stay; /* user asked to keep this PCON on via UI */ char needed; /* transiently set by deps analysis */ @@ -139,6 +134,7 @@ struct sai_power { lws_sorted_usec_list_t sul_idle; lws_sorted_usec_list_t sul_pcon_check; /* periodic check for cold start */ + lws_sorted_usec_list_t sul_monitor; /* periodic energy monitoring */ const char *power_off; @@ -187,6 +183,10 @@ void saip_set_stay(const char *pcon_name, int stay_on); int saip_queue_stay_info(saip_server_t *sps); + +int +saip_queue_energy_report(saip_server_t *sps); + saip_pcon_t * saip_pcon_by_name(struct sai_power *power, const char *name); diff --git a/src/power/p-sai.c b/src/power/p-sai.c index 2ee1855..2afc42a 100644 --- a/src/power/p-sai.c +++ b/src/power/p-sai.c @@ -250,10 +250,56 @@ sul_pcon_check_cb(lws_sorted_usec_list_t *sul) } void +sul_broadcast_energy_cb(lws_sorted_usec_list_t *sul) +{ + int polled = 0; + + /* 1. Trigger monitoring on all Tasmota PCONs */ + 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->ss_tasmota_monitor) { + polled++; + /* Reset RX position for new response */ + pc->monitor_rx_pos = 0; + /* SS request triggers the HTTP GET */ + if (lws_ss_request_tx(pc->ss_tasmota_monitor)) + lwsl_warn("%s: Failed to trigger monitor request for %s\n", __func__, pc->name); + else + lwsl_notice("%s: Triggered polling for %s\n", __func__, pc->name); + } else { + lwsl_warn("%s: PCON %s has no monitor SS\n", __func__, pc->name); + } + + } lws_end_foreach_dll(p); + + if (!polled) + lwsl_notice("%s: No PCONs polled\n", __func__); + + /* 2. Queue energy report to server (sends whatever latest data we have) */ + /* We iterate servers, though usually only one */ + lws_start_foreach_dll_safe(struct lws_dll2 *, mp, mp1, + power.sai_server_owner.head) { + saip_server_t *sps = lws_container_of(mp, struct saip_server, list); + int queued = saip_queue_energy_report(sps); + if (queued) { + lwsl_notice("%s: Queued energy report for server\n", __func__); + if (lws_ss_request_tx(sps->ss)) /* Request write to send the report */ + lwsl_warn("%s: Failed to trigger monitor request\n", __func__); + } + } lws_end_foreach_dll_safe(mp, mp1); + + /* Schedule next check (e.g., every 5 seconds) */ + lws_sul_schedule(power.context, 0, &power.sul_monitor, sul_broadcast_energy_cb, 5 * LWS_US_PER_SEC); +} + +void saip_pcon_start_check(void) { /* Trigger immediate check */ lws_sul_schedule(power.context, 0, &power.sul_pcon_check, sul_pcon_check_cb, 1); + /* Trigger immediate energy monitor check */ + lws_sul_schedule(power.context, 0, &power.sul_monitor, sul_broadcast_energy_cb, 1); } static int diff --git a/src/power/p-smartplug.c b/src/power/p-smartplug.c index 02d1f21..58a496e 100644 --- a/src/power/p-smartplug.c +++ b/src/power/p-smartplug.c @@ -32,11 +32,67 @@ LWS_SS_USER_TYPEDEF } saip_smartplug_t; static lws_ss_state_return_t +saip_spc_rx(void *userobj, const uint8_t *buf, size_t len, int flags) +{ + saip_smartplug_t *pss = (saip_smartplug_t *)userobj; + saip_pcon_t *pc = (saip_pcon_t *)lws_ss_opaque_from_user(pss); + tasmota_parse_t tp; + struct lws_ss_handle *h = lws_ss_from_user(pss); + + /* We only care about rx on the monitor channel */ + if (h != pc->ss_tasmota_monitor) + return 0; + + if (pc->monitor_rx_pos + len > sizeof(pc->monitor_rx_buf) - 1) { + lwsl_warn("%s: monitor rx buffer overflow\n", __func__); + pc->monitor_rx_pos = 0; + return 0; + } + + memcpy(pc->monitor_rx_buf + pc->monitor_rx_pos, buf, len); + pc->monitor_rx_pos += len; + + if (flags & LWSSS_FLAG_EOM) { + pc->monitor_rx_buf[pc->monitor_rx_pos] = '\0'; + + // lwsl_notice("%s: monitor rx: %s\n", __func__, pc->monitor_rx_buf); + + memset(&tp, 0, sizeof(tp)); + lws_tokenize_init(&tp.ts, pc->monitor_rx_buf, + LWS_TOKENIZE_F_NO_FLOATS | + LWS_TOKENIZE_F_MINUS_NONTERM); + + int parse_ret = saip_parse_tasmota_status(&tp); + if (parse_ret == 1) { + /* Success, update latest data and timestamp */ + pc->latest_data = tp.td; + pc->last_monitor_time = lws_now_usecs(); + lwsl_notice("%s: Parsed Tasmota data for %s: %u W, %u V\n", + __func__, pc->name, pc->latest_data.active_power_w, pc->latest_data.voltage_v); + } else + lwsl_warn("%s: Failed to parse Tasmota data for %s (ret %d)\n", __func__, pc->name, parse_ret); + + + pc->monitor_rx_pos = 0; + } + + return 0; +} + +static lws_ss_state_return_t saip_spc_state(void *userobj, void *sh, lws_ss_constate_t state, lws_ss_tx_ordinal_t ack) { saip_smartplug_t *pss = (saip_smartplug_t *)userobj; - const char *op_url = (const char *)lws_ss_opaque_from_user(pss); + saip_pcon_t *pc = (saip_pcon_t *)lws_ss_opaque_from_user(pss); + const char *op_url = NULL; + + if (lws_ss_from_user(pss) == pc->ss_tasmota_on) + op_url = pc->url_on; + else if (lws_ss_from_user(pss) == pc->ss_tasmota_off) + op_url = pc->url_off; + else if (lws_ss_from_user(pss) == pc->ss_tasmota_monitor) + op_url = pc->url_monitor; // lwsl_user("%s: %s, ord 0x%x\n", __func__, lws_ss_state_name((int)state), // (unsigned int)ack); @@ -44,13 +100,25 @@ saip_spc_state(void *userobj, void *sh, lws_ss_constate_t state, switch (state) { case LWSSSCS_CREATING: - lwsl_notice("%s: binding ss to %s\n", __func__, op_url); + if (!op_url) { + lwsl_err("%s: creating unknown SS\n", __func__); + return LWSSSSRET_DISCONNECT_ME; + } + // lwsl_notice("%s: binding ss to %s\n", __func__, op_url); if (lws_ss_set_metadata(lws_ss_from_user(pss), "url", op_url, strlen(op_url))) lwsl_warn("%s: unable to set metadata\n", __func__); break; + case LWSSSCS_CONNECTED: + /* If this is a monitor connection, we might want to reset the buffer */ + if (lws_ss_from_user(pss) == pc->ss_tasmota_monitor) { + pc->monitor_rx_pos = 0; + /* Note: SS with "GET" method will automatically send request */ + } + return lws_ss_request_tx(lws_ss_from_user(pss)); + default: break; } @@ -60,6 +128,7 @@ saip_spc_state(void *userobj, void *sh, lws_ss_constate_t state, LWS_SS_INFO("sai_power_smartplug", saip_smartplug_t) .state = saip_spc_state, + .rx = saip_spc_rx, }; void @@ -74,7 +143,7 @@ saip_ss_create_tasmota() lws_snprintf(pc->url_on, sizeof(pc->url_on), "%s/cm?cmnd=Power%%20On", pc->url); if (lws_ss_create(power.context, 0, &ssi_saip_smartplug_t, - (void *)pc->url_on, + (void *)pc, &pc->ss_tasmota_on, NULL, NULL)) lwsl_err("%s: %s: failed to create ON smartplug secure stream %s\n", __func__, pc->name, pc->url_on); @@ -82,18 +151,20 @@ saip_ss_create_tasmota() lws_snprintf(pc->url_off, sizeof(pc->url_off), "%s/cm?cmnd=Power%%20Off", pc->url); if (lws_ss_create(power.context, 0, &ssi_saip_smartplug_t, - (void *)pc->url_off, + (void *)pc, &pc->ss_tasmota_off, NULL, NULL)) lwsl_err("%s: %s: failed to create OFF smartplug secure stream %s\n", __func__, pc->name, pc->url_off); lws_snprintf(pc->url_monitor, sizeof(pc->url_monitor), - "%s?m=1", pc->url); + "%s/?m=1", pc->url); if (lws_ss_create(power.context, 0, &ssi_saip_smartplug_t, - (void *)pc->url_monitor, + (void *)pc, &pc->ss_tasmota_monitor, NULL, NULL)) lwsl_err("%s: %s: failed to create MONITOR smartplug secure stream %s\n", __func__, pc->name, pc->url_monitor); + else + lwsl_ss_warn(pc->ss_tasmota_monitor, "============================ creating monitor SS for %s", pc->url_monitor); } } lws_end_foreach_dll(px); diff --git a/src/power/p-tasmota-monitor.c b/src/power/p-tasmota-monitor.c index 8d81e1b..e4036c2 100644 --- a/src/power/p-tasmota-monitor.c +++ b/src/power/p-tasmota-monitor.c @@ -103,8 +103,10 @@ saip_parse_tasmota_status(tasmota_parse_t *tp) } } - if (n == LWS_ARRAY_SIZE(tokens)) + if (n == LWS_ARRAY_SIZE(tokens)) { + // lwsl_notice("%s: unknown token '%.*s'\n", __func__, (int)tp->ts.token_len, tp->ts.token); continue; + } break; case LWS_TOKZE_INTEGER: diff --git a/src/power/p-ws-server.c b/src/power/p-ws-server.c index c3c605a..238bec7 100644 --- a/src/power/p-ws-server.c +++ b/src/power/p-ws-server.c @@ -36,6 +36,61 @@ static const lws_struct_map_t lsm_schema_power_state[] = { "com.warmcat.sai.powerstate"), }; +/* + * (Structs and maps removed - now in common/include/private.h and common/struct-metadata.c) + */ + +int +saip_queue_energy_report(saip_server_t *sps) +{ + sai_pcon_energy_report_t report; + struct lwsac *ac = NULL; + saip_server_link_t *m; + int r = 0; + int count = 0; + + if (!sps->ss) + return 0; + + m = (saip_server_link_t *)lws_ss_to_user_object(sps->ss); + + memset(&report, 0, sizeof(report)); + + 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); + + /* Only include if we have valid data (checked within last 60s?) */ + if (pc->last_monitor_time && + lws_now_usecs() - pc->last_monitor_time < 60 * LWS_US_PER_SEC) { + sai_pcon_energy_report_item_t *item = + lwsac_use_zero(&ac, sizeof(*item), 1024); + + if (item) { + lws_strncpy(item->name, pc->name, sizeof(item->name)); + item->data = pc->latest_data; + lws_dll2_add_tail(&item->list, &report.items); + count++; + } + } else { + if (pc->last_monitor_time) + lwsl_notice("%s: Stale monitor data for %s (age %llus)\n", __func__, pc->name, (unsigned long long)(lws_now_usecs() - pc->last_monitor_time) / LWS_US_PER_SEC); + else + lwsl_notice("%s: No monitor data for %s\n", __func__, pc->name); + } + } lws_end_foreach_dll(p); + + if (count) { + lwsl_notice("%s: Queuing energy report with %d items\n", __func__, count); + r = sai_ss_serialize_queue_helper(sps->ss, &m->bl_pwr_to_srv, + lsm_schema_pcon_energy, + LWS_ARRAY_SIZE(lsm_schema_pcon_energy), + &report); + } + + lwsac_free(&ac); + return r; +} + void saip_notify_server_power_state(const char *pcon_name, int up, int down) { @@ -139,49 +194,72 @@ saip_m_rx(void *userobj, const uint8_t *buf, size_t len, int flags) struct lejp_ctx ctx; lwsl_notice("%s: len %d, flags: %d (saip_server_t %p)\n", __func__, (int)len, flags, (void *)sps); + lwsl_hexdump_notice(buf, len); /* lwsl_hexdump_notice(buf, len); */ memset(&a, 0, sizeof(a)); a.map_st[0] = lsm_schema_stay; a.map_entries_st[0] = LWS_ARRAY_SIZE(lsm_schema_stay); + a.map_st[1] = lsm_schema_pcon_control; + a.map_entries_st[1] = LWS_ARRAY_SIZE(lsm_schema_pcon_control); a.ac_block_size = 512; lws_struct_json_init_parse(&ctx, NULL, &a); if (lejp_parse(&ctx, (uint8_t *)buf, (int)len) >= 0 && a.dest) { - sai_stay_t *stay = (sai_stay_t *)a.dest; - - // {"schema":"com.warmcat.sai.power.stay","builder_name":"ubuntu_rpi4","stay_on":1} - - lwsl_warn("%s: received stay %s: %d\n", __func__, stay->builder_name, stay->stay_on); - - /* - * We received a stay request for a builder. - * We need to find which PCON controls this builder and update its manual_stay. - * But wait, 'sai_stay_t' is typically per-builder. - * We should map this back to the PCON. - */ - - 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); - lws_start_foreach_dll(struct lws_dll2 *, b_node, pc->registered_builders_owner.head) { - saip_builder_t *sb = lws_container_of(b_node, saip_builder_t, list); - if (!strcmp(sb->name, stay->builder_name)) { - lwsl_notice("%s: Mapping stay for builder '%s' to PCON '%s'\n", - __func__, sb->name, pc->name); - /* Update PCON stay state */ - pc->manual_stay = stay->stay_on; - - /* If stay is cleared, schedule power off check */ - if (!stay->stay_on) - saip_pcon_start_check(); - else { - /* If stay is set, ensure it is on immediately */ - saip_switch(pc, 1); - } - goto found; + + if (a.top_schema_index == 1) { + /* PCON Control */ + sai_pcon_control_t *ctl = (sai_pcon_control_t *)a.dest; + saip_pcon_t *pc = saip_pcon_by_name(&power, ctl->pcon_name); + + if (pc) { + lwsl_notice("%s: PCON Control '%s' -> %d\n", __func__, pc->name, ctl->on); + pc->manual_stay = ctl->on; + if (ctl->on) { + saip_switch(pc, 1); + } else { + saip_pcon_start_check(); } - } lws_end_foreach_dll(b_node); - } lws_end_foreach_dll(p); + } else { + lwsl_warn("%s: Unknown PCON '%s'\n", __func__, ctl->pcon_name); + } + } else { + /* Stay */ + sai_stay_t *stay = (sai_stay_t *)a.dest; + + // {"schema":"com.warmcat.sai.power.stay","builder_name":"ubuntu_rpi4","stay_on":1} + + lwsl_warn("%s: received stay %s: %d\n", __func__, stay->builder_name, stay->stay_on); + + /* + * We received a stay request for a builder. + * We need to find which PCON controls this builder and update its manual_stay. + * But wait, 'sai_stay_t' is typically per-builder. + * We should map this back to the PCON. + */ + + 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); + lws_start_foreach_dll(struct lws_dll2 *, b_node, pc->registered_builders_owner.head) { + saip_builder_t *sb = lws_container_of(b_node, saip_builder_t, list); + if (!strcmp(sb->name, stay->builder_name)) { + lwsl_notice("%s: Mapping stay for builder '%s' to PCON '%s'\n", + __func__, sb->name, pc->name); + /* Update PCON stay state */ + pc->manual_stay = stay->stay_on; + + /* If stay is cleared, schedule power off check */ + if (!stay->stay_on) + saip_pcon_start_check(); + else { + /* If stay is set, ensure it is on immediately */ + saip_switch(pc, 1); + } + goto found; + } + } lws_end_foreach_dll(b_node); + } lws_end_foreach_dll(p); + } found: lwsac_free(&a.ac); diff --git a/src/server/s-power.c b/src/server/s-power.c index 132eacb..01b25ac 100644 --- a/src/server/s-power.c +++ b/src/server/s-power.c @@ -37,12 +37,15 @@ #include "s-private.h" +/* + * (Structs and maps removed - now in common/include/private.h and common/struct-metadata.c) + */ + int sais_power_rx(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t bl, unsigned int ss_flags) { - struct lejp_ctx ctx; - lws_struct_args_t a; + lws_struct_args_t *a = &pss->a; sai_power_state_t *ps; const lws_struct_map_t lsm_schema_map_power[] = { LSM_SCHEMA(sai_power_state_t, NULL, lsm_power_state, @@ -53,26 +56,52 @@ sais_power_rx(struct vhd *vhd, struct pss *pss, uint8_t *buf, LSM_SCHEMA(sai_stay_state_update_t, NULL, lsm_stay_state_update, "com.warmcat.sai.stay_state_update"), + /* We just passthrough PCON energy reports to the web side */ + LSM_SCHEMA(sai_pcon_energy_report_t, NULL, /* Use correct struct/map */ + lsm_pcon_energy_report, + "com.warmcat.sai.pcon_energy"), }; + int n; /* This is a message from sai-power */ - lwsl_notice("RX from sai-power: %.*s\n", (int)bl, (const char *)buf); - - memset(&a, 0, sizeof(a)); - a.map_st[0] = lsm_schema_map_power; - a.map_entries_st[0] = LWS_ARRAY_SIZE(lsm_schema_map_power); - a.ac_block_size = 512; - - lws_struct_json_init_parse(&ctx, NULL, &a); - if (lejp_parse(&ctx, buf, (int)bl) < 0 || !a.dest) { - lwsl_warn("Failed to parse msg from sai-power\n"); - lwsac_free(&a.ac); - return 1; + /* lwsl_notice("RX from sai-power: %.*s\n", (int)bl, (const char *)buf); */ + + if (ss_flags & LWSSS_FLAG_SOM) { + memset(a, 0, sizeof(*a)); + a->top_schema_index = -1; + a->map_st[0] = lsm_schema_map_power; + a->map_entries_st[0] = LWS_ARRAY_SIZE(lsm_schema_map_power); + a->ac_block_size = 512; + + lws_struct_json_init_parse(&pss->ctx_power, NULL, a); + lws_buflist_destroy_all_segments(&pss->power_rx_cache); + } + + /* We always cache the fragment until we know what it is */ + if (lws_buflist_append_segment(&pss->power_rx_cache, buf, bl) < 0) { + lwsl_err("%s: failed to append to power_rx_cache\n", __func__); + return -1; + } + + n = lejp_parse(&pss->ctx_power, buf, (int)bl); + if (n < 0 && n != LEJP_CONTINUE) { + lwsl_warn("Failed to parse msg from sai-power %s (schema idx %d)\n", + lejp_error_to_string(n), a->top_schema_index); + goto bail; } - switch (a.top_schema_index) { + if (a->top_schema_index == 3) /* com.warmcat.sai.pcon_energy */ + goto passthru; + + if (n == LEJP_CONTINUE) + return 0; + + if (!a->dest) + goto bail; + + switch (a->top_schema_index) { case 0: /* powerstate */ - ps = (sai_power_state_t *)a.dest; + ps = (sai_power_state_t *)a->dest; lwsl_notice("%s: powerstate received: %d %d\n", __func__, ps->powering_up, ps->powering_down); if (ps->powering_up) { lwsl_notice("sai-power is powering up: %s\n", ps->host); @@ -86,7 +115,7 @@ sais_power_rx(struct vhd *vhd, struct pss *pss, uint8_t *buf, break; case 1: { - sai_power_managed_builders_t *pmb = (sai_power_managed_builders_t *)a.dest; + sai_power_managed_builders_t *pmb = (sai_power_managed_builders_t *)a->dest; char q[256]; /* @@ -134,7 +163,7 @@ sais_power_rx(struct vhd *vhd, struct pss *pss, uint8_t *buf, break; } case 2: { - sai_stay_state_update_t *ssu = (sai_stay_state_update_t *)a.dest; + sai_stay_state_update_t *ssu = (sai_stay_state_update_t *)a->dest; sai_plat_t *sp; lwsl_notice("%s: Received stay_state_update for %s, stay_on=%d\n", @@ -158,16 +187,190 @@ sais_power_rx(struct vhd *vhd, struct pss *pss, uint8_t *buf, break; } + case 3: /* com.warmcat.sai.pcon_energy */ +passthru: + /* + * This is an energy report from sai-power. + * We want to broadcast it to all connected web interfaces (sai-web). + * We have been buffering in pss->power_rx_cache until we identified the schema. + * Now we forward whatever is in the cache (which includes the current 'buf'). + */ + { + lws_wsmsg_info_t info; + uint8_t *p, *lin; + size_t tlen = lws_buflist_total_len(&pss->power_rx_cache); + + /* + * We need to linearize the cache to send it out with LWS_PRE. + */ + p = malloc(LWS_PRE + tlen); + if (!p) { + lwsl_err("%s: OOM forwarding energy report\n", __func__); + lws_buflist_destroy_all_segments(&pss->power_rx_cache); + break; + } + + lin = p + LWS_PRE; + size_t copied = 0; + + /* Drain the buflist into our linear buffer */ + while (copied < tlen) { + uint8_t *seg; + size_t slen; + + slen = lws_buflist_next_segment_len(&pss->power_rx_cache, &seg); + if (!slen) + break; + + memcpy(lin + copied, seg, slen); + copied += slen; + lws_buflist_use_segment(&pss->power_rx_cache, slen); + } + + memset(&info, 0, sizeof(info)); + info.private_source_idx = SAI_WEBSRV_PB__GENERATED; + info.buf = p + LWS_PRE; + info.len = tlen; + info.ss_flags = LWSSS_FLAG_SOM; /* We always send what we have as a start */ + + if (ss_flags & LWSSS_FLAG_EOM) + info.ss_flags |= LWSSS_FLAG_EOM; + + /* + * If we are continuing (n == LEJP_CONTINUE), we flushed the buffer + * so subsequent calls will append new data to empty buflist and flush it immediately. + * However, sais_websrv_broadcast expects SOM/EOM to be correct for the whole message. + * + * If we buffered the START of the message, we set SOM. + * If the incoming chunk was EOM, we set EOM. + * + * What if we have intermediate chunks? + * + * If we are in passthru, we cleared the cache above. + * + * Wait, if we are in passthru state, we shouldn't re-set SOM for every chunk. + * We need to track if we already sent SOM. + * + * But here we only enter `passthru` if `top_schema_index` matches. + * This happens for the FIRST chunk (once schema matches) AND subsequent chunks. + * + * Problem: `top_schema_index` remains 3 for subsequent chunks. + * + * So we need to know if we are flushing the FIRST part (SOM) or a later part. + * + * `ss_flags & LWSSS_FLAG_SOM` tells us if the CURRENT chunk was the start of the message. + * + * If `ss_flags & SOM`, then `info.ss_flags |= SOM`. + * If `ss_flags & EOM`, then `info.ss_flags |= EOM`. + * + * This seems correct because we are effectively delaying the processing. + * If we buffered chunks 1 and 2, and now processing chunk 2 (which made schema valid), + * chunk 1 had SOM. `pss->power_rx_cache` contains chunk 1 + chunk 2. + * So the aggregate buffer DOES start with SOM content. + * + * If we are processing chunk 3 (schema already known), we append to cache, then flush. + * Cache contains just chunk 3. + * Chunk 3 does NOT have SOM. + * So we shouldn't set SOM. + * + * BUT `ss_flags` belongs to the current `buf` (chunk 3). + * If `ss_flags` has SOM, then our buffer starts with SOM. + * + * So `info.ss_flags = ss_flags` is ALMOST correct, except that we might have accumulated + * previous chunks which HAD SOM, even if the current chunk doesn't. + * + * If `fragment_cache` was non-empty before we appended `buf`, then we are continuing a buffer. + * Wait, we appended `buf` to `cache` at the top of the function. + * + * If `ss_flags` has SOM, then the cache definitely starts with SOM. + * + * If `ss_flags` does NOT have SOM, but we have older data in cache? + * That older data MUST be the start of the message (because we flush on schema detection). + * + * Wait, if we flush on schema detection, we flush the START. + * Subsequent chunks will be appended to EMPTY cache, then flushed. + * + * So if cache has data, and we are flushing... + * + * Case 1: First chunk(s). Schema found. `ss_flags` might be SOM (if single chunk) or NOT (if 2nd chunk). + * If 2nd chunk triggers match, `ss_flags` is !SOM. But cache contains Chunk 1 (SOM) + Chunk 2. + * So we must set SOM if the *cache* contains the start. + * + * We can track `pss->power_rx_cache_had_som`. + * Or we can just rely on `a->top_schema_index == -1` -> we are at start. + * Once matched, we are flushing the start. + * + * Actually, simpler: + * We only buffer if we DON'T know the schema. + * + * If we know the schema (3), we are in passthru mode. + * + * If we just transitioned to schema 3 (match occurred in this chunk), we flush everything. This flush INCLUDES the start. So send SOM. + * + * If we were ALREADY in schema 3 (subsequent chunks), we just forward `buf`. + * + * But wait, my logic "always append to cache" means `buf` is in cache. + * + * If I flush cache every time `passthru` is hit: + * - First time (match): Cache has Start + ... + Current. Flush. Send SOM. + * - Next time: Cache has Next Chunk. Flush. Send !SOM. + * + * How do I know if it's the "First time"? + * `a->top_schema_index` is persistent in `pss`. + * + * Valid point: `lws_struct` parser state persists. + * + * Issue: `lejp` doesn't tell me "I just matched schema". + * + * But I can check if `pss->power_rx_cache` contains more than `bl`. + * If `total_len > bl`, then we have buffered data -> We are sending the start -> SOM. + * + * Exception: What if `buf` is the FIRST chunk and it matched? + * `total_len == bl`. But `ss_flags` has SOM. + * + * So logic: + * `info.ss_flags = (ss_flags & LWSSS_FLAG_EOM);` + * `if (total_len > bl || (ss_flags & LWSSS_FLAG_SOM)) info.ss_flags |= LWSSS_FLAG_SOM;` + * + * This handles: + * - Single chunk (SOM+EOM): len==bl, flags=SOM. Result: SOM+EOM. + * - Multi chunk, 1st (match): len==bl, flags=SOM. Result: SOM. + * - Multi chunk, 2nd (match): len > bl, flags=!SOM. Result: SOM. (Correct, as it contains start). + * - Multi chunk, 3rd (already matched): + * Wait, if already matched, we still append and flush? + * If we flush every time, cache is empty between calls. + * So for 3rd chunk, `total_len == bl`. `flags`=!SOM. Result: !SOM. Correct. + * + */ + + if (tlen > bl || (ss_flags & LWSSS_FLAG_SOM)) + info.ss_flags |= LWSSS_FLAG_SOM; + else + info.ss_flags &= (unsigned int)~LWSSS_FLAG_SOM; + + /* Broadcast to all websrv connections (i.e. all sai-web instances) */ + sais_websrv_broadcast_REQUIRES_LWS_PRE(vhd->h_ss_websrv, &info); + + free(p); + } + break; default: lwsl_warn("%s: unknown schema\n", __func__); + /* If it's not our schema, we must clear the cache so it doesn't leak into next message */ + lws_buflist_destroy_all_segments(&pss->power_rx_cache); break; } - lwsac_free(&a.ac); +bail: + lwsac_free(&a->ac); return 0; } +/* + * (Structs and maps removed - now in common/include/private.h and common/struct-metadata.c) + */ + int sais_power_tx(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t bl) { @@ -210,6 +413,39 @@ sais_power_tx(struct vhd *vhd, struct pss *pss, uint8_t *buf, size_t bl) return 0; } + if (pss->pcon_control_owner.head) { + /* + * Pending PCON control message to send to power + */ + sai_pcon_control_t *s = lws_container_of(pss->pcon_control_owner.head, + sai_pcon_control_t, list); + lws_struct_serialize_t *js; + + js = lws_struct_json_serialize_create(lsm_schema_pcon_control, + LWS_ARRAY_SIZE(lsm_schema_pcon_control), 0, s); + if (!js) { + lwsl_warn("%s: failed to serialize pcon control\n", __func__); + return 1; + } + + n = (int)lws_struct_json_serialize(js, p, lws_ptr_diff_size_t(end, p), &w); + lws_struct_json_serialize_destroy(&js); + + lwsl_wsi_notice(pss->wsi, "%s: server issuing pcon control notice\n", __func__); + sai_dump_stderr(start, w); + + lws_dll2_remove(&s->list); + free(s); + + flags = lws_write_ws_flags(LWS_WRITE_TEXT, 1, 1); + + if (lws_write(pss->wsi, start, w, flags) < 0) + return -1; + + lws_callback_on_writable(pss->wsi); + return 0; + } + n = 0; lws_start_foreach_dll(struct lws_dll2 *, px, vhd->pending_plats.head) { sais_plat_t *pl = lws_container_of(px, sais_plat_t, list); diff --git a/src/server/s-private.h b/src/server/s-private.h index 779b12e..f9c9178 100644 --- a/src/server/s-private.h +++ b/src/server/s-private.h @@ -134,6 +134,7 @@ struct pss { struct lws_dll2 same; /* owner: vhd.builders */ struct lws_buflist *onward_reassembly; + struct lws_buflist *power_rx_cache; sqlite3 *pdb_artifact; sqlite3_blob *blob_artifact; @@ -142,6 +143,7 @@ struct pss { lws_dll2_owner_t task_cancel_owner; /* sai_platform_t builder offers */ lws_dll2_owner_t rebuild_owner; lws_dll2_owner_t stay_owner; + lws_dll2_owner_t pcon_control_owner; lws_dll2_owner_t aft_owner; /* for statefully spooling artifact info */ lws_dll2_owner_t res_owner; /* sai_resource_requisition_t * owner of resource objects related @@ -151,6 +153,7 @@ struct pss { * messages to builder */ lws_dll2_owner_t viewer_state_owner; lws_struct_args_t a; + struct lejp_ctx ctx_power; const char *server_name; diff --git a/src/server/s-ws-web.c b/src/server/s-ws-web.c index 267eb28..5cdfde7 100644 --- a/src/server/s-ws-web.c +++ b/src/server/s-ws-web.c @@ -58,6 +58,10 @@ static lws_struct_map_t lsm_browser_taskreset[] = { LSM_CARRAY (sai_browse_rx_evinfo_t, event_hash, "uuid"), }; +/* + * (Structs and maps removed - now in common/include/private.h and common/struct-metadata.c) + */ + static lws_struct_map_t lsm_browser_platreset[] = { LSM_CARRAY (sai_browse_rx_platreset_t, event_uuid, "event_uuid"), LSM_CARRAY (sai_browse_rx_platreset_t, platform, "platform"), @@ -93,6 +97,8 @@ static const lws_struct_map_t lsm_schema_json_map[] = { "com.warmcat.sai.platreset"), LSM_SCHEMA (sai_stay_t, NULL, lsm_stay, "com.warmcat.sai.stay"), + LSM_SCHEMA (sai_pcon_control_t, NULL, lsm_pcon_control, + /* shares struct */ "com.warmcat.sai.pcon_control"), LSM_SCHEMA (sai_browse_rx_taskinfo_t, NULL, lsm_browser_taskinfo, "com.warmcat.sai.taskinfo") }; @@ -107,6 +113,7 @@ enum { SAIS_WS_WEBSRV_RX_REBUILD, SAIS_WS_WEBSRV_RX_PLATRESET, SAIS_WS_WEBSRV_RX_STAY, + SAIS_WS_WEBSRV_RX_PCON_CONTROL, SAIS_WS_WEBSRV_RX_TASKINFO, }; @@ -402,8 +409,8 @@ websrvss_ws_rx(void *userobj, const uint8_t *buf, size_t len, int flags) sai_db_result_t r; int n; - // lwsl_user("%s: len %d, flags: %d\n", __func__, (int)len, flags); - // lwsl_hexdump_info(buf, len); + lwsl_user("%s: len %d, flags: %d\n", __func__, (int)len, flags); + lwsl_hexdump_info(buf, len); memset(&a, 0, sizeof(a)); a.map_st[0] = lsm_schema_json_map; @@ -470,6 +477,30 @@ websrvss_ws_rx(void *userobj, const uint8_t *buf, size_t len, int flags) lwsac_free(&a.ac); break; } + case SAIS_WS_WEBSRV_RX_PCON_CONTROL: + { + sai_pcon_control_t *ctl = (sai_pcon_control_t *)a.dest; + + lwsl_notice("%s: pcon control received from web: %s -> %d\n", + __func__, ctl->pcon_name, ctl->on); + + lws_start_foreach_dll(struct lws_dll2 *, p, + m->vhd->sai_powers.head) { + struct pss *pss_power = lws_container_of(p, struct pss, same); + sai_pcon_control_t *s; + + s = malloc(sizeof(*s)); + if (s) { + *s = *ctl; + lws_dll2_add_tail(&s->list, &pss_power->pcon_control_owner); + lws_callback_on_writable(pss_power->wsi); + lwsl_wsi_notice(pss_power->wsi, "queued pcon control on power conn"); + } + } lws_end_foreach_dll(p); + + lwsac_free(&a.ac); + break; + } case SAIS_WS_WEBSRV_RX_EVENTDELETE: ei = (sai_browse_rx_evinfo_t *)a.dest; diff --git a/src/web/w-ws-browser.c b/src/web/w-ws-browser.c index 8570e0e..2d13c49 100644 --- a/src/web/w-ws-browser.c +++ b/src/web/w-ws-browser.c @@ -38,6 +38,10 @@ * For decoding specific event data request from browser */ +/* + * (Structs and maps removed - now in common/include/private.h and common/struct-metadata.c) + */ + static lws_struct_map_t lsm_browser_evinfo[] = { LSM_CARRAY (sai_browse_rx_evinfo_t, event_hash, "event_hash"), }; @@ -86,6 +90,8 @@ static const lws_struct_map_t lsm_schema_json_map_bwsrx[] = { "com.warmcat.sai.platreset"), LSM_SCHEMA (sai_stay_t, NULL, lsm_stay, "com.warmcat.sai.stay"), + LSM_SCHEMA (sai_pcon_control_t, NULL, lsm_pcon_control, + /* shares struct */ "com.warmcat.sai.pcon_control"), }; enum { @@ -100,6 +106,7 @@ enum { SAIM_WS_BROWSER_RX_REBUILD, SAIM_WS_BROWSER_RX_PLATRESET, SAIM_WS_BROWSER_RX_STAY, + SAIM_WS_BROWSER_RX_PCON_CONTROL, }; @@ -571,6 +578,9 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, sai_cancel_t *can; int m, ret = -1; + lwsl_notice("%s: len %d, flags: %d\n", __func__, (int)bl, ss_flags); + /* lwsl_hexdump_notice(buf, bl); */ + memset(&a, 0, sizeof(a)); /* * pss->js_api_version defaults to 1 (from ESTABLISHED callback). @@ -669,6 +679,17 @@ saiw_ws_json_rx_browser(struct vhd *vhd, struct pss *pss, uint8_t *buf, */ break; + case SAIM_WS_BROWSER_RX_PCON_CONTROL: + if (!sais_conn_auth(pss)) { + lwsl_err("%s: pcon control didn't like auth\n", __func__); + goto auth_error; + } + lwsl_notice("%s: web: received pcon control req\n", __func__); + + /* Forward to sai-server via websrv link */ + /* We rely on the fallthrough to queue the message buffer to websrv */ + break; + case SAIM_WS_BROWSER_RX_TASKREBUILDLASTSTEP: if (!sais_conn_auth(pss)) goto auth_error; diff --git a/src/web/w-ws-server.c b/src/web/w-ws-server.c index baaeaa7..a1d3e82 100644 --- a/src/web/w-ws-server.c +++ b/src/web/w-ws-server.c @@ -36,6 +36,10 @@ static lws_struct_map_t lsm_websrv_evinfo[] = { LSM_CARRAY (sai_browse_rx_evinfo_t, event_hash, "event_hash"), }; +/* + * (Structs and maps removed - now in common/include/private.h and common/struct-metadata.c) + */ + const lws_struct_map_t lsm_schema_json_map[] = { LSM_SCHEMA (sai_browse_rx_evinfo_t, NULL, lsm_websrv_evinfo, /* shares struct */ "sai-taskchange"), @@ -55,6 +59,8 @@ const lws_struct_map_t lsm_schema_json_map[] = { LSM_SCHEMA(sai_power_managed_builders_t, NULL, lsm_power_managed_builders_list, "com.warmcat.sai.power_managed_builders"), + LSM_SCHEMA (sai_pcon_energy_report_t, NULL, lsm_pcon_energy_report, + /* shares struct */ "com.warmcat.sai.pcon_energy"), }; enum { @@ -67,6 +73,7 @@ enum { SAIS_WS_WEBSRV_RX_TASKACTIVITY, SAIS_WS_WEBSRV_RX_BUILD_METRIC, SAIS_WS_WEBSRV_RX_POWER_MANAGED_BUILDERS, + SAIS_WS_WEBSRV_RX_PCON_ENERGY, }; /* @@ -122,6 +129,7 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags) case SAIS_WS_WEBSRV_RX_TASKACTIVITY: case SAIS_WS_WEBSRV_RX_SAI_BUILDERS: case SAIS_WS_WEBSRV_RX_POWER_MANAGED_BUILDERS: + case SAIS_WS_WEBSRV_RX_PCON_ENERGY: saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, buf, len, lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, @@ -139,6 +147,7 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags) case SAIS_WS_WEBSRV_RX_EVENTCHANGE: case SAIS_WS_WEBSRV_RX_SAI_BUILDERS: case SAIS_WS_WEBSRV_RX_POWER_MANAGED_BUILDERS: + case SAIS_WS_WEBSRV_RX_PCON_ENERGY: saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, buf, len, lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, @@ -267,6 +276,10 @@ saiw_lp_rx(void *userobj, const uint8_t *buf, size_t len, int flags) saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, buf, len - (unsigned int)n, lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM)); break; + case SAIS_WS_WEBSRV_RX_PCON_ENERGY: + saiw_ws_broadcast_browsers_REQUIRES_LWS_PRE(vhd, buf, len - (unsigned int)n, + lws_write_ws_flags(LWS_WRITE_TEXT, flags & LWSSS_FLAG_SOM, flags & LWSSS_FLAG_EOM)); + break; } cleanup_parse_allocs:
Page fetched 0s ago, creation time: 9ms (vhost etag hits: 0%, cache hits: 0%)