hi team i am trying to download the file from the server directly to the flash iam continusly getting this error <err> node_fota_dl: >>> Node firmware download_client ERROR: -113, iam using the nrf5340 as the app core and for gsm external modem actualy the firmware is recving the two firmware from the server the first firmware will download in mcuboot but the other firmware which iam trying to download that should be dwnloaded n the flash directly but iam continusly getting the error . I have checked there is no issue with the modm everything is fine the parsing is fine when iam copying the url which iam printhing t when iam copying that url that is not able to download . but when iam hitting the api get firmware api in the postman that is being able to downloa
#include <zephyr/net/mqtt.h>
#include <zephyr/net/socket.h>
#include <zephyr/logging/log.h>
#include <zephyr/data/json.h>
#include <modem/modem_info.h>
#include <zephyr/sys/reboot.h>
#include <zephyr/sys/util.h>
#include "aws_iot/aws-certs.h"
#include "iot_connections.h"
#include <stdio.h>
#include <stdlib.h>
#include <ctype.h>
#include "dms_handler.h"
#include <net/fota_download.h>
#include <dfu/dfu_target_mcuboot.h>
#include <zephyr/net/sntp.h>
#include <zephyr/drivers/flash.h>
#include <dfu/dfu_target.h>
#include "time.h"
#include "node_fota_download.h"
#if defined(CONFIG_AWS_FOTA)
#include <net/aws_fota.h>
#include <zephyr/net/net_ip.h>
#endif
#if defined(CONFIG_AWS_IOT_PROVISION_CERTIFICATES)
#include CONFIG_AWS_IOT_CERTIFICATES_FILE
#endif
#define MY_DEVICE_IMEI CONFIG_DEVICE_IMEI
atomic_t fota_active = ATOMIC_INIT(0);
LOG_MODULE_REGISTER(aws, LOG_LEVEL_DBG);
#define SNTP_SERVER "0.pool.ntp.org"
#define SEC_TAG_DMS 10
#define SEC_TAG_CATTLE 11
#if !defined(CONFIG_AWS_IOT_BROKER_HOST_NAME_APP)
BUILD_ASSERT(sizeof(CONFIG_AWS_IOT_BROKER_HOST_NAME) > 1, "AWS IoT hostname not set");
BUILD_ASSERT(CONFIG_AWS_IOT_BROKER_HOST_NAME_MAX_LEN >= sizeof(CONFIG_AWS_IOT_BROKER_HOST_NAME) - 1,"AWS IoT host name static buffer too small ""Increase CONFIG_AWS_IOT_BROKER_HOST_NAME_MAX_LEN");
#endif
#if defined(CONFIG_AWS_IOT_IPV6)
#define AWS_AF_FAMILY AF_INET6
#else
#define AWS_AF_FAMILY AF_INET
#endif
static char *client_id_buf = CONFIG_AWS_THING_NAME;
#define AWS_IOT_SHADOW_REQUEST_STRING ""
static char dms_host_name_buf[] ="a2sdhgohxx19hy-ats.iot.ap-south-1.amazonaws.com";
static char cattle_host_name_buf[] ="vanix-mybovin-iot-core.voltrackvanix.com";
#define MQTT_PORT 8883
static char aws_host_name_buf[CONFIG_AWS_IOT_BROKER_HOST_NAME_MAX_LEN + 1];
static struct aws_iot_app_topic_data app_topic_data;
static struct mqtt_client client;
static struct sockaddr_storage broker;
static struct mqtt_client client_dms;
K_MUTEX_DEFINE(dms_client_mutex); /* protects client_dms socket lifecycle across threads */
bool fota_check_done = false;
static struct sockaddr_storage broker_dms;
static uint8_t rx_buffer_dms[1024];
static uint8_t tx_buffer_dms[1024];
static char client_id_buf_dms[40];
static bool sntp_synced = false;
static char rx_buffer[CONFIG_AWS_IOT_MQTT_RX_TX_BUFFER_LEN];
static char tx_buffer[CONFIG_AWS_IOT_MQTT_RX_TX_BUFFER_LEN];
static char payload_buf[CONFIG_AWS_IOT_MQTT_PAYLOAD_BUFFER_LEN];
static aws_iot_evt_handler_t module_evt_handler;
static atomic_t disconnect_requested;
static atomic_t connection_poll_active;
atomic_t aws_iot_disconnected = ATOMIC_INIT(1);
#define SEC_TAG_S3 12
#define DMS_RECONNECT_MAX_ATTEMPTS 10
#define DMS_RECONNECT_RETRY_DELAY_SEC 5
#define DMS_RECONNECT_STACK_SIZE 4096
#define DMS_RECONNECT_PRIORITY K_PRIO_PREEMPT(7)
#define DMS_RECONNECT_CYCLE_BACKOFF_SEC (5 * 60)
#define DMS_RECONNECT_MAX_CYCLE_FAILURES 6
void dms_trigger_reconnect(void);
int dms_iot_send(const struct aws_iot_data *const tx_data)
{
if (client_dms.transport.tls.sock < 0) {
LOG_WRN("DMS socket not connected, dropping publish");
log_to_flash("DMS socket not connected, dropping publish");
dms_trigger_reconnect();
return -ENOTCONN;
}
struct mqtt_publish_param param = {0};
param.message.topic.qos = tx_data->qos;
param.message.topic.topic.utf8 = tx_data->topic.str;
param.message.topic.topic.size = tx_data->topic.len;
param.message.payload.data = tx_data->ptr;
param.message.payload.len = tx_data->len;
param.dup_flag = 0;
param.retain_flag = 0;
param.message_id = k_cycle_get_32();
return mqtt_publish(&client_dms, ¶m);
}
static K_SEM_DEFINE(connection_poll_sem, 0, 1);
static int connect_error_translate(const int err)
{
switch (err)
{
case 0:
return AWS_IOT_CONNECT_RES_SUCCESS;
case -ECHILD:
return AWS_IOT_CONNECT_RES_ERR_NETWORK;
case -EACCES:
return AWS_IOT_CONNECT_RES_ERR_NOT_INITD;
case -ENOEXEC:
return AWS_IOT_CONNECT_RES_ERR_BACKEND;
case -EINVAL:
return AWS_IOT_CONNECT_RES_ERR_PRV_KEY;
case -EOPNOTSUPP:
return AWS_IOT_CONNECT_RES_ERR_CERT;
case -ECONNREFUSED:
return AWS_IOT_CONNECT_RES_ERR_CERT_MISC;
case -ETIMEDOUT:
return AWS_IOT_CONNECT_RES_ERR_TIMEOUT_NO_DATA;
case -ENOMEM:
return AWS_IOT_CONNECT_RES_ERR_NO_MEM;
case -EINPROGRESS:
return AWS_IOT_CONNECT_RES_ERR_ALREADY_CONNECTED;
default:
LOG_ERR("AWS broker connect failed %d", err);
return AWS_IOT_CONNECT_RES_ERR_MISC;
}
}
static void aws_iot_notify_event(const struct aws_iot_evt *aws_iot_evt)
{
if ((module_evt_handler != NULL) && (aws_iot_evt != NULL))
{
module_evt_handler(aws_iot_evt);
}
else
{
LOG_ERR("Library event handler not registered, or empty event");
}
}
#define DMS_SERVER_IP "65.1.96.216"
#define DMS_SERVER_PORT 3600
#define AUTH_TOKEN "a12b4e1bed6a2aba55081ff587e670c15a795526669f8abf416d2eb785562c82"
static char fota_presigned_url[2560] = {0};
static char fota_latest_version[32] = {0};
static char fota_response[6144];
static char fota_host_buf[256];
static char fota_path_buf[2560];
static void fota_retry_work_fn(struct k_work *work)
{
LOG_INF("FOTA retry: reconnecting to fetch fresh presigned URL...");
fota_check_done = false;
k_sem_give(&connection_poll_sem);
}
static K_WORK_DELAYABLE_DEFINE(fota_retry_work, fota_retry_work_fn);
/* Minimal completion handler for the raw-download test — just logs
* the result and releases the fota_active lock so normal MQTT/publish
* resumes. Does NOT trigger a BLE push (that comes later, once this
* basic "does the file land in flash" step is confirmed working). */
static void raw_download_done_cb(bool success, const char *node_mac,
off_t flash_addr, size_t image_len,
const uint8_t sha256[32])
{
ARG_UNUSED(node_mac);
ARG_UNUSED(sha256);
if (success) {
LOG_INF("=== RAW DOWNLOAD TEST: SUCCESS, %u bytes written at 0x%lx ===",
(unsigned)image_len, (long)flash_addr);
log_to_flash("Raw download test: success");
fota_report_status(MY_DEVICE_IMEI, fota_latest_version, "Success");
} else {
LOG_ERR("=== RAW DOWNLOAD TEST: FAILED ===");
log_to_flash("Raw download test: failed");
fota_report_status(MY_DEVICE_IMEI, fota_latest_version, "Failed");
}
atomic_set(&fota_active, 0);
}
void fota_dl_handler(const struct fota_download_evt *evt)
{
switch (evt->id) {
case FOTA_DOWNLOAD_EVT_PROGRESS:
LOG_INF("Download progress to flash: %d%%", evt->progress);
break;
case FOTA_DOWNLOAD_EVT_FINISHED:
LOG_INF("SUCCESS: File streamed and saved to flash memory!");
log_to_flash("Download to flash finished successfully");
fota_report_status(MY_DEVICE_IMEI, fota_latest_version, "Success");
/* Direct download complete — release lock without writing MCUboot headers or rebooting */
atomic_set(&fota_active, 0);
break;
case FOTA_DOWNLOAD_EVT_ERROR:
LOG_ERR("FOTA download error — scheduling retry in 15s");
log_to_flash("FOTA download error — scheduling retry in 15s");
atomic_set(&fota_active, 0);
fota_report_status(MY_DEVICE_IMEI, fota_latest_version, "Failed");
fota_check_done = false;
k_work_reschedule(&fota_retry_work, K_SECONDS(15));
break;
default:
break;
}
}
int fota_report_status(const char *device_imei, const char *version,
const char *status)
{
int sock, err;
struct sockaddr_in server;
char body[256];
char request[640];
char response[512];
snprintf(body, sizeof(body), "{\"message\":3,\"device_id\":\"%s\"," "\"installed_version\":\"%s\"," "\"status\":\"%s\"," "\"timestamp\":%lld}", device_imei, version, status, (long long)(k_uptime_get() / 1000));
snprintf(request, sizeof(request),
"POST /v1/api/device/firmware/update-upload-status HTTP/1.1\r\n"
"Host: %s:%d\r\n"
"Authorization: %s\r\n"
"Content-Type: application/json\r\n"
"Content-Length: %d\r\n"
"Connection: close\r\n\r\n%s",
DMS_SERVER_IP, DMS_SERVER_PORT, AUTH_TOKEN,
strlen(body), body);
sock = socket(AF_INET, SOCK_STREAM, IPPROTO_TCP);
if (sock < 0) { LOG_ERR("fota_report socket failed"); return -errno; }
server.sin_family = AF_INET;
server.sin_port = htons(DMS_SERVER_PORT);
net_addr_pton(AF_INET, DMS_SERVER_IP, &server.sin_addr);
err = connect(sock, (struct sockaddr *)&server, sizeof(server));
if (err < 0) { close(sock); return -errno; }
send(sock, request, strlen(request), 0);
recv(sock, response, sizeof(response) - 1, 0);
LOG_INF("FOTA status report: %s", status);
close(sock);
return 0;
}
#define DEVICE_REG_MAGIC 0xD00D2026UL
#define DEVICE_REG_PRODUCT_TYPE 1
#define DEVICE_REG_PRODUCT_CATEGORY_ID 5
#define DEVICE_REG_HARDWARE_TYPE 1
#define DEVICE_REG_HARDWARE_VERSION 1
#define DEVICE_REG_CLIENT_ID 11
#define DEVICE_REG_NAME_PREFIX "VANIX-GW"
#ifndef CONFIG_FLASH_DEVICE_REG_ADD
#define CONFIG_FLASH_DEVICE_REG_ADD (CONFIG_FLASH_FOTA_VERSION_ADD + 0x1000)
#endif
struct device_reg_flag {
uint32_t magic;
char imei[24];
};
static bool device_registration_done(void)
{
const struct device *flash_dev = DEVICE_DT_GET(DT_ALIAS(spi_flash0));
struct device_reg_flag flag = {0};
if (!device_is_ready(flash_dev)) {
LOG_ERR("Device-reg: flash not ready");
return false;
}
if (flash_read(flash_dev, CONFIG_FLASH_DEVICE_REG_ADD, &flag, sizeof(flag)) != 0) {
return false;
}
if (flag.magic == DEVICE_REG_MAGIC &&
strncmp(flag.imei, MY_DEVICE_IMEI, sizeof(flag.imei) - 1) == 0) {
return true;
}
return false;
}
static void device_registration_mark_done(void)
{
const struct device *flash_dev = DEVICE_DT_GET(DT_ALIAS(spi_flash0));
struct device_reg_flag flag = {0};
if (!device_is_ready(flash_dev)) {
LOG_ERR("Device-reg: flash not ready, cannot persist flag");
return;
}
flag.magic = DEVICE_REG_MAGIC;
strncpy(flag.imei, MY_DEVICE_IMEI, sizeof(flag.imei) - 1);
flash_erase(flash_dev, CONFIG_FLASH_DEVICE_REG_ADD, 4096);
flash_write(flash_dev, CONFIG_FLASH_DEVICE_REG_ADD, &flag, sizeof(flag));
LOG_INF("Device-reg: registration flag persisted for IMEI %s", MY_DEVICE_IMEI);
log_to_flash("Device registration flag persisted");
}
void device_registration_reset(void)
{
const struct device *flash_dev = DEVICE_DT_GET(DT_ALIAS(spi_flash0));
if (!device_is_ready(flash_dev)) {
LOG_ERR("Device-reg: flash not ready, cannot reset flag");
return;
}
flash_erase(flash_dev, CONFIG_FLASH_DEVICE_REG_ADD, 4096);
LOG_INF("Device-reg: registration flag cleared, will re-register on next connect");
log_to_flash("Device registration flag cleared");
}
int device_register_with_dms(void)
{
int sock, err;
struct sockaddr_in server;
char body[384];
char request[768];
char response[512];
if (device_registration_done()) {
LOG_DBG("Device-reg: already registered, skipping");
return 0;
}
snprintf(body, sizeof(body),
"{\"product_type\":%d,\"product_category_id\":%d,"
"\"hardware_type\":%d,\"device_sim_number\":\"%s\","
"\"hardware_version\":%d,\"device_name\":\"%s-%s\","
"\"client_id\":%d,\"device_imei\":\"%s\"}",
DEVICE_REG_PRODUCT_TYPE,
DEVICE_REG_PRODUCT_CATEGORY_ID,
DEVICE_REG_HARDWARE_TYPE,
MY_DEVICE_IMEI,
DEVICE_REG_HARDWARE_VERSION,
DEVICE_REG_NAME_PREFIX, MY_DEVICE_IMEI,
DEVICE_REG_CLIENT_ID,
MY_DEVICE_IMEI);
snprintf(request, sizeof(request),
"POST /v1/api/device/add-device HTTP/1.1\r\n"
"Host: %s:%d\r\n"
"Authorization: %s\r\n"
"Content-Type: application/json\r\n"
"Content-Length: %d\r\n"
"Connection: close\r\n\r\n%s",
DMS_SERVER_IP, DMS_SERVER_PORT, AUTH_TOKEN,
strlen(body), body);
sock = socket(AF_INET, SOCK_STREAM, IPPROTO_TCP);
if (sock < 0) {
LOG_ERR("Device-reg: socket failed");
return -errno;
}
struct timeval tv = { .tv_sec = 15, .tv_usec = 0 };
setsockopt(sock, SOL_SOCKET, SO_SNDTIMEO, &tv, sizeof(tv));
setsockopt(sock, SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof(tv));
server.sin_family = AF_INET;
server.sin_port = htons(DMS_SERVER_PORT);
net_addr_pton(AF_INET, DMS_SERVER_IP, &server.sin_addr);
err = connect(sock, (struct sockaddr *)&server, sizeof(server));
if (err < 0) {
LOG_ERR("Device-reg: connect failed, errno %d", errno);
close(sock);
return -errno;
}
send(sock, request, strlen(request), 0);
memset(response, 0, sizeof(response));
recv(sock, response, sizeof(response) - 1, 0);
close(sock);
LOG_INF("Device-reg response: %s", response);
if (strstr(response, "\"success\":true") != NULL ||
strstr(response, "\"success\": true") != NULL ||
strstr(response, "\"status\":201") != NULL ||
strstr(response, "\"status\": 201") != NULL) {
LOG_INF("Device-reg: registration successful");
log_to_flash("Device registered to DMS successfully");
device_registration_mark_done();
return 0;
}
LOG_WRN("Device-reg: registration failed (or device already exists on server)");
log_to_flash("Device registration to DMS failed");
return -EIO;
}
int fota_check_and_start(const char *device_imei, const char *target_node_mac)
{
int sock, err;
struct sockaddr_in server;
char url_path[128];
char request[512];
int total = 0, received;
static int fota_attempt_count = 0;
LOG_INF("Entering fota_check_and_start (Direct Flash Stream Mode)");
log_to_flash("FOTA: entering fota_check_and_start");
atomic_set(&fota_active, 1);
k_sleep(K_MSEC(300));
if (fota_attempt_count >= 3) {
LOG_ERR("FOTA failed 3 times, giving up");
atomic_set(&fota_active, 0);
fota_attempt_count = 0;
return 1;
}
snprintf(url_path, sizeof(url_path),
"/v1/api/device/firmware/get-firmware-by-device-id/%s",
device_imei);
snprintf(request, sizeof(request),
"GET %s HTTP/1.1\r\n"
"Host: %s:%d\r\n"
"Authorization: %s\r\n"
"User-Agent: BMSDevice\r\n"
"Connection: close\r\n\r\n",
url_path,
DMS_SERVER_IP,
DMS_SERVER_PORT,
AUTH_TOKEN);
sock = socket(AF_INET, SOCK_STREAM, IPPROTO_TCP);
if (sock < 0) { LOG_ERR("fota_check socket failed"); return -errno; }
struct timeval tv = { .tv_sec = 15, .tv_usec = 0 };
setsockopt(sock, SOL_SOCKET, SO_SNDTIMEO, &tv, sizeof(tv));
setsockopt(sock, SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof(tv));
server.sin_family = AF_INET;
server.sin_port = htons(DMS_SERVER_PORT);
net_addr_pton(AF_INET, DMS_SERVER_IP, &server.sin_addr);
err = connect(sock, (struct sockaddr *)&server, sizeof(server));
if (err < 0) { close(sock); return -errno; }
send(sock, request, strlen(request), 0);
memset(fota_response, 0, sizeof(fota_response));
total = 0;
struct pollfd pfd = { .fd = sock, .events = POLLIN };
while (1) {
int p = poll(&pfd, 1, 10000);
if (p <= 0) break;
received = recv(sock, fota_response + total,
sizeof(fota_response) - total - 1, 0);
if (received <= 0) break;
total += received;
}
close(sock);
LOG_INF("FOTA response total bytes: %d", total);
char *json_body = strstr(fota_response, "\r\n\r\n");
if (!json_body) {
LOG_ERR("No HTTP body in DMS response");
atomic_set(&fota_active, 0);
return -EINVAL;
}
char *url_start_ptr = strstr(fota_response, "\"presigned_url\":\"");
if (!url_start_ptr) {
url_start_ptr = strstr(fota_response, "\"presigned_url\" : \"");
if (!url_start_ptr) {
LOG_ERR("No presigned_url found");
atomic_set(&fota_active, 0);
return -EINVAL;
}
url_start_ptr += strlen("\"presigned_url\" : \"");
} else {
url_start_ptr += strlen("\"presigned_url\":\"");
}
char *url_end_ptr = url_start_ptr;
while (*url_end_ptr) {
if (*url_end_ptr == '"' && *(url_end_ptr - 1) != '\\') break;
url_end_ptr++;
}
if (*url_end_ptr != '"') {
LOG_ERR("Could not find end of presigned_url");
atomic_set(&fota_active, 0);
return -EINVAL;
}
size_t url_len = url_end_ptr - url_start_ptr;
memset(fota_presigned_url, 0, sizeof(fota_presigned_url));
memcpy(fota_presigned_url, url_start_ptr,
MIN(url_len, sizeof(fota_presigned_url) - 1));
static char unescape_buf[2560];
{
memset(unescape_buf, 0, sizeof(unescape_buf));
char *src = fota_presigned_url;
char *dst = unescape_buf;
while (*src && (dst - unescape_buf) < (int)sizeof(unescape_buf) - 1) {
if (*src == '\\' && *(src + 1) == '/') {
*dst++ = '/'; src += 2;
} else if (*src == '\\' && *(src + 1) == '"') {
*dst++ = '"'; src += 2;
} else if (*src == '\\' && *(src + 1) == 'u' &&
isxdigit((unsigned char)src[2]) && isxdigit((unsigned char)src[3]) &&
isxdigit((unsigned char)src[4]) && isxdigit((unsigned char)src[5])) {
/* \uXXXX escape — common for '&' (\u0026) in presigned URL
* query strings from some JSON serializers. */
char hex[5] = { src[2], src[3], src[4], src[5], '\0' };
unsigned int codepoint = (unsigned int)strtoul(hex, NULL, 16);
if (codepoint < 0x80) {
*dst++ = (char)codepoint;
}
src += 6;
} else {
*dst++ = *src++;
}
}
*dst = '\0';
strncpy(fota_presigned_url, unescape_buf, sizeof(fota_presigned_url) - 1);
fota_presigned_url[sizeof(fota_presigned_url) - 1] = '\0';
}
char *const fota_host = fota_host_buf;
char *const fota_path = fota_path_buf;
char *host_start = strstr(fota_presigned_url, "://");
if (!host_start) {
LOG_ERR("Invalid URL");
atomic_set(&fota_active, 0);
return -EINVAL;
}
host_start += 3;
char *path_start = strchr(host_start, '/');
if (!path_start) {
LOG_ERR("Invalid URL path");
atomic_set(&fota_active, 0);
return -EINVAL;
}
memset(fota_host, 0, 256);
memset(fota_path, 0, sizeof(fota_path_buf));
strncpy(fota_host, host_start, MIN((size_t)(path_start - host_start), 256 - 1));
strncpy(fota_path, path_start, sizeof(fota_path_buf) - 1);
/* ── DEBUG: confirm the host/path split is correct ── */
LOG_INF("=== PARSED: host=[%s] ===", fota_host);
LOG_INF("=== PARSED: path_len=%d ===", (int)strlen(fota_path));
LOG_INF("Disconnecting MQTT prior to direct flash streaming...");
/* NOTE: previously this captured dms_sock/cattle_sock *before* calling
* mqtt_disconnect(), then close()'d that stale value afterward.
* mqtt_disconnect() already closes the underlying socket internally,
* so that stale close() was a double-close on an fd number that could
* be reused by ANY other thread's socket in the meantime (SNTP, epoch
* anchor, RSSI poll all run concurrently — visible in the logs right
* around this teardown). Closing a reused fd corrupts whichever
* unrelated socket now owns that number, which is a very plausible
* explanation for "brand new S3 connect fails with EHOSTUNREACH right
* after this teardown, even though DNS just succeeded and the same
* URL works fine from Postman." Fix: always read the CURRENT socket
* value right before closing it, never a value captured earlier. */
k_mutex_lock(&dms_client_mutex, K_FOREVER);
mqtt_disconnect(&client_dms);
k_sleep(K_MSEC(300));
if (client_dms.transport.tls.sock >= 0) {
close(client_dms.transport.tls.sock);
}
client_dms.transport.tls.sock = -1;
k_mutex_unlock(&dms_client_mutex);
atomic_set(&disconnect_requested, 1);
mqtt_disconnect(&client);
k_sleep(K_MSEC(200));
if (client.transport.tls.sock >= 0) {
close(client.transport.tls.sock);
}
client.transport.tls.sock = -1;
atomic_set(&aws_iot_disconnected, 1);
k_sleep(K_SECONDS(5));
tls_credential_delete(SEC_TAG_S3, TLS_CREDENTIAL_CA_CERTIFICATE);
tls_credential_add(SEC_TAG_S3, TLS_CREDENTIAL_CA_CERTIFICATE, amazon_root_ca_1, amazon_root_ca1_len);
const char *fota_path_no_slash = fota_path;
if (fota_path_no_slash[0] == '/') {
fota_path_no_slash++;
}
err = node_fota_download_from_host_path(
target_node_mac,
fota_host,
fota_path_no_slash,
raw_download_done_cb);
if (err) {
LOG_ERR("Download start failed: %d", err);
atomic_set(&fota_active, 0);
fota_report_status(device_imei, fota_latest_version, "Failed");
fota_attempt_count++;
return err;
}
LOG_INF("File download directly to flash initialized...");
log_to_flash("File download directly to flash initialized...");
return 0;
}
#if defined(CONFIG_AWS_FOTA)
static void aws_fota_cb_handler(struct aws_fota_event *fota_evt)
{
struct aws_iot_evt aws_iot_evt = {0};
if (fota_evt == NULL)
{
return;
}
switch (fota_evt->id)
{
case AWS_FOTA_EVT_START:
LOG_DBG("AWS_FOTA_EVT_START");
aws_iot_evt.type = AWS_IOT_EVT_FOTA_START;
break;
case AWS_FOTA_EVT_DONE:
LOG_DBG("AWS_FOTA_EVT_DONE");
fota_report_status(MY_DEVICE_IMEI, fota_latest_version, "Success");
aws_iot_evt.type = AWS_IOT_EVT_FOTA_DONE;
aws_iot_evt.data.image = fota_evt->image;
break;
case AWS_FOTA_EVT_ERASE_PENDING:
LOG_DBG("AWS_FOTA_EVT_ERASE_PENDING");
aws_iot_evt.type = AWS_IOT_EVT_FOTA_ERASE_PENDING;
break;
case AWS_FOTA_EVT_ERASE_DONE:
LOG_DBG("AWS_FOTA_EVT_ERASE_DONE");
aws_iot_evt.type = AWS_IOT_EVT_FOTA_ERASE_DONE;
break;
case AWS_FOTA_EVT_ERROR:
LOG_ERR("AWS_FOTA_EVT_ERROR");
log_to_flash("AWS_FOTA_EVT_ERROR");
fota_report_status(MY_DEVICE_IMEI, fota_latest_version, "Failed");
aws_iot_evt.type = AWS_IOT_EVT_FOTA_ERROR;
break;
case AWS_FOTA_EVT_DL_PROGRESS:
LOG_DBG("AWS_FOTA_EVT_DL_PROGRESS, (%d%%)",
fota_evt->dl.progress);
aws_iot_evt.type = AWS_IOT_EVT_FOTA_DL_PROGRESS;
aws_iot_evt.data.fota_progress = fota_evt->dl.progress;
break;
default:
LOG_ERR("Unknown FOTA event");
return;
}
aws_iot_notify_event(&aws_iot_evt);
}
#endif
static int aws_iot_set_host_name(char *const host_name, size_t host_name_len)
{
if (host_name == NULL)
{
return -EINVAL;
}
if (host_name_len >= sizeof(aws_host_name_buf))
{
LOG_ERR("Host name too long: %d bytes (max %d)", host_name_len, sizeof(aws_host_name_buf) - 1);
return -ENOMEM;
}
memcpy(aws_host_name_buf, host_name, host_name_len);
aws_host_name_buf[host_name_len] = '\0';
return 0;
}
static int publish_get_payload(struct mqtt_client *const c, size_t length)
{
if (length > sizeof(payload_buf))
{
LOG_ERR("Incoming MQTT message too large for payload buffer");
return -EMSGSIZE;
}
return mqtt_readall_publish_payload(c, payload_buf, length);
}
static K_SEM_DEFINE(cloud_ready_sem, 0, 1);
void aws_iot_wait_for_connection(void)
{
k_sem_take(&cloud_ready_sem, K_FOREVER);
}
static void mqtt_evt_handler(struct mqtt_client *const c, const struct mqtt_evt *mqtt_evt)
{
int err;
struct aws_iot_evt aws_iot_evt = {0};
if (c == &client_dms) {
switch (mqtt_evt->type) {
case MQTT_EVT_CONNACK:
if (mqtt_evt->param.connack.return_code == 0) {
LOG_INF("DMS MQTT client connected!");
log_to_flash("DMS MQTT client connected!");
if (app_topic_data.list_count > 0) {
struct mqtt_subscription_list sub_list = {
.list = app_topic_data.list,
.list_count = app_topic_data.list_count,
.message_id = k_cycle_get_32()
};
err = mqtt_subscribe(c, &sub_list);
if (err) {
LOG_ERR("DMS Subscribe failed: %d", err);
} else {
LOG_INF("DMS subscribed to %d topics", app_topic_data.list_count);
}
}
} else {
LOG_ERR("DMS CONNACK error: %d", mqtt_evt->param.connack.return_code);
}
break;
case MQTT_EVT_PUBLISH: {
const struct mqtt_publish_param *p = &mqtt_evt->param.publish;
err = publish_get_payload(c, p->message.payload.len);
if (err == 0) {
dms_route_message(
p->message.topic.topic.utf8,
p->message.topic.topic.size,
(const uint8_t *)payload_buf,
p->message.payload.len);
}
if (p->message.topic.qos == MQTT_QOS_1_AT_LEAST_ONCE) {
const struct mqtt_puback_param ack = { .message_id = p->message_id };
mqtt_publish_qos1_ack(c, &ack);
}
break;
}
case MQTT_EVT_PUBACK:
break;
case MQTT_EVT_PINGRESP:
break;
case MQTT_EVT_DISCONNECT:
LOG_WRN("DMS MQTT client disconnected: %d", mqtt_evt->result);
log_to_flash("DMS MQTT client disconnected");
break;
default:
break;
}
return;
}
switch (mqtt_evt->type) {
case MQTT_EVT_CONNACK:
if (mqtt_evt->param.connack.return_code) {
aws_iot_evt.data.err = mqtt_evt->param.connack.return_code;
aws_iot_evt.type = AWS_IOT_EVT_ERROR;
aws_iot_notify_event(&aws_iot_evt);
break;
}
LOG_DBG("Cattle MQTT client connected!");
log_to_flash("Cattle MQTT client connected!");
if (app_topic_data.list_count > 0) {
struct mqtt_subscription_list sub_list = {
.list = app_topic_data.list,
.list_count = app_topic_data.list_count,
.message_id = k_cycle_get_32(),
};
err = mqtt_subscribe(c, &sub_list);
if (!err) {
LOG_INF("Cattle subscribing to %d topics...", app_topic_data.list_count);
log_to_flash("Cattle subscribed");
}
}
aws_iot_evt.data.persistent_session =
!IS_ENABLED(CONFIG_MQTT_CLEAN_SESSION) &&
mqtt_evt->param.connack.session_present_flag;
aws_iot_evt.type = AWS_IOT_EVT_CONNECTED;
aws_iot_notify_event(&aws_iot_evt);
aws_iot_evt.type = AWS_IOT_EVT_READY;
aws_iot_notify_event(&aws_iot_evt);
k_sem_give(&cloud_ready_sem);
break;
case MQTT_EVT_DISCONNECT:
aws_iot_evt.data.err = AWS_IOT_DISCONNECT_MISC;
if (atomic_get(&disconnect_requested)) {
aws_iot_evt.data.err = AWS_IOT_DISCONNECT_USER_REQUEST;
}
atomic_set(&aws_iot_disconnected, 1);
aws_iot_evt.type = AWS_IOT_EVT_DISCONNECTED;
aws_iot_notify_event(&aws_iot_evt);
break;
case MQTT_EVT_PUBLISH: {
const struct mqtt_publish_param *p = &mqtt_evt->param.publish;
err = publish_get_payload(c, p->message.payload.len);
if (err) break;
if (p->message.topic.qos == MQTT_QOS_1_AT_LEAST_ONCE) {
const struct mqtt_puback_param ack = { .message_id = p->message_id };
mqtt_publish_qos1_ack(c, &ack);
}
char topic_cmp[128] = {0};
memcpy(topic_cmp, p->message.topic.topic.utf8, MIN(p->message.topic.topic.size, sizeof(topic_cmp) - 1));
if (strstr(topic_cmp, "get-log-file") != NULL) {
handle_log_file_request((const uint8_t *)payload_buf, p->message.payload.len);
} else if (strstr(topic_cmp, "restart-device") != NULL) {
handle_device_restart_request((const uint8_t *)payload_buf, p->message.payload.len);
} else if (strstr(topic_cmp, "check-heart-beat-device") != NULL) {
handle_heartbeat_request((const uint8_t *)payload_buf, p->message.payload.len);
} else {
LOG_WRN("Received message on unhandled topic: %s", topic_cmp);
log_to_flash("Received message on unhandled topic");
}
aws_iot_evt.type = AWS_IOT_EVT_DATA_RECEIVED;
aws_iot_evt.data.msg.ptr = payload_buf;
aws_iot_evt.data.msg.len = p->message.payload.len;
aws_iot_evt.data.msg.topic.str = p->message.topic.topic.utf8;
aws_iot_evt.data.msg.topic.len = p->message.topic.topic.size;
aws_iot_notify_event(&aws_iot_evt);
break;
}
case MQTT_EVT_PUBACK:
aws_iot_evt.type = AWS_IOT_EVT_PUBACK;
aws_iot_evt.data.message_id = mqtt_evt->param.puback.message_id;
aws_iot_notify_event(&aws_iot_evt);
break;
case MQTT_EVT_PINGRESP:
aws_iot_evt.type = AWS_IOT_EVT_PINGRESP;
aws_iot_notify_event(&aws_iot_evt);
break;
default:
break;
}
}
static int broker_init_dms(void)
{
int err;
struct addrinfo *result;
struct addrinfo *addr;
struct addrinfo hints = {
.ai_family = AF_INET,
.ai_socktype = SOCK_STREAM,
};
err = getaddrinfo(dms_host_name_buf, NULL, &hints, &result);
if (err) {
LOG_ERR("getaddrinfo DMS, error %d", err);
return -ECHILD;
}
addr = result;
while (addr != NULL) {
if (addr->ai_addrlen == sizeof(struct sockaddr_in)) {
struct sockaddr_in *broker4 = (struct sockaddr_in *)&broker_dms;
broker4->sin_addr.s_addr = ((struct sockaddr_in *)addr->ai_addr)->sin_addr.s_addr;
broker4->sin_family = AF_INET;
broker4->sin_port = htons(MQTT_PORT);
break;
}
addr = addr->ai_next;
}
freeaddrinfo(result);
return 0;
}
static int certificates_provision_dms(void)
{
static bool dms_certs_added;
if (dms_certs_added) return 0;
tls_credential_delete(SEC_TAG_DMS, TLS_CREDENTIAL_CA_CERTIFICATE);
tls_credential_delete(SEC_TAG_DMS, TLS_CREDENTIAL_PRIVATE_KEY);
tls_credential_delete(SEC_TAG_DMS, TLS_CREDENTIAL_SERVER_CERTIFICATE);
int err = tls_credential_add(SEC_TAG_DMS, TLS_CREDENTIAL_CA_CERTIFICATE, dms_ca_certificate, sizeof(dms_ca_certificate));
if (err < 0) return err;
err = tls_credential_add(SEC_TAG_DMS, TLS_CREDENTIAL_PRIVATE_KEY, dms_private_key, sizeof(dms_private_key));
if (err < 0) return err;
err = tls_credential_add(SEC_TAG_DMS, TLS_CREDENTIAL_SERVER_CERTIFICATE, dms_device_certificate, sizeof(dms_device_certificate));
if (err < 0) return err;
dms_certs_added = true;
LOG_INF("DMS Certs OK");
log_to_flash("DMS Certs OK");
return 0;
}
static int certificates_provision_cattle(void)
{
static bool cattle_certs_added;
if (cattle_certs_added) return 0;
tls_credential_delete(SEC_TAG_CATTLE, TLS_CREDENTIAL_CA_CERTIFICATE);
tls_credential_delete(SEC_TAG_CATTLE, TLS_CREDENTIAL_PRIVATE_KEY);
tls_credential_delete(SEC_TAG_CATTLE, TLS_CREDENTIAL_SERVER_CERTIFICATE);
int err = tls_credential_add(SEC_TAG_CATTLE, TLS_CREDENTIAL_CA_CERTIFICATE, cattle_ca_certificate, sizeof(cattle_ca_certificate));
if (err < 0) return err;
err = tls_credential_add(SEC_TAG_CATTLE, TLS_CREDENTIAL_PRIVATE_KEY, cattle_private_key, sizeof(cattle_private_key));
if (err < 0) return err;
err = tls_credential_add(SEC_TAG_CATTLE, TLS_CREDENTIAL_SERVER_CERTIFICATE, cattle_device_certificate, sizeof(cattle_device_certificate));
if (err < 0) return err;
cattle_certs_added = true;
LOG_INF("Cattle Certs OK");
log_to_flash("Cattle Certs OK");
return 0;
}
static int client_broker_init_dms(struct mqtt_client *const client)
{
mqtt_client_init(client);
client->keepalive = 120;
if (broker_init_dms() != 0) return -ECHILD;
client->broker = &broker_dms;
client->evt_cb = mqtt_evt_handler;
client->client_id.utf8 = client_id_buf_dms;
client->client_id.size = strlen(client_id_buf_dms);
client->rx_buf = rx_buffer_dms;
client->rx_buf_size = sizeof(rx_buffer_dms);
client->tx_buf = tx_buffer_dms;
client->tx_buf_size = sizeof(tx_buffer_dms);
client->transport.type = MQTT_TRANSPORT_SECURE;
static sec_tag_t tags_dms[] = { SEC_TAG_DMS };
client->transport.tls.config.sec_tag_list = tags_dms;
client->transport.tls.config.sec_tag_count = 1;
client->transport.tls.config.peer_verify = 2;
client->transport.tls.config.hostname = dms_host_name_buf;
return certificates_provision_dms();
}
int connect_client_dms(void)
{
int err;
k_mutex_lock(&dms_client_mutex, K_FOREVER);
err = client_broker_init_dms(&client_dms);
if (err) {
LOG_ERR("client_broker_init_dms, error: %d", err);
k_mutex_unlock(&dms_client_mutex);
return err;
}
err = mqtt_connect(&client_dms);
k_mutex_unlock(&dms_client_mutex);
if (err) {
LOG_ERR("mqtt_connect DMS, error: %d", err);
return err;
}
return 0;
}
static atomic_t dms_reconnect_in_progress = ATOMIC_INIT(0);
static atomic_t dms_reconnect_workq_started = ATOMIC_INIT(0);
static struct k_work_q dms_reconnect_work_q;
K_THREAD_STACK_DEFINE(dms_reconnect_stack_area, DMS_RECONNECT_STACK_SIZE);
static int dms_reconnect_cycle_failures = 0;
static void dms_reconnect_work_fn(struct k_work *work);
static K_WORK_DELAYABLE_DEFINE(dms_reconnect_work, dms_reconnect_work_fn);
static void dms_reconnect_work_fn(struct k_work *work)
{
extern bool is_network_connected;
bool dms_recovered = false;
if (atomic_set(&dms_reconnect_in_progress, 1) == 1) {
return;
}
LOG_WRN("DMS reconnect: starting recovery (max %d attempts)", DMS_RECONNECT_MAX_ATTEMPTS);
log_to_flash("DMS reconnect: starting recovery");
for (int attempt = 1; attempt <= DMS_RECONNECT_MAX_ATTEMPTS; attempt++) {
if (atomic_get(&fota_active)) {
LOG_INF("DMS reconnect: FOTA active, standing down (not a failure)");
log_to_flash("DMS reconnect: standing down for FOTA");
atomic_set(&dms_reconnect_in_progress, 0);
return;
}
if (atomic_get(&aws_iot_disconnected) == 1) {
LOG_WRN("DMS reconnect: Cattle/AWS disconnected mid-cycle, standing down");
log_to_flash("DMS reconnect: aborted, Cattle disconnected mid-cycle");
break;
}
if (client_dms.transport.tls.sock >= 0) {
dms_recovered = true;
break;
}
if (is_network_connected) {
int err = connect_client_dms();
if (err == 0 && client_dms.transport.tls.sock >= 0) {
dms_recovered = true;
break;
}
k_mutex_lock(&dms_client_mutex, K_FOREVER);
if (client_dms.transport.tls.sock >= 0) {
close(client_dms.transport.tls.sock);
client_dms.transport.tls.sock = -1;
}
k_mutex_unlock(&dms_client_mutex);
}
if (attempt < DMS_RECONNECT_MAX_ATTEMPTS) {
k_sleep(K_SECONDS(DMS_RECONNECT_RETRY_DELAY_SEC));
}
}
if (dms_recovered) {
dms_reconnect_cycle_failures = 0;
} else {
dms_reconnect_cycle_failures++;
if (atomic_get(&aws_iot_disconnected) == 1) {
k_sleep(K_SECONDS(2));
sys_reboot(SYS_REBOOT_COLD);
} else if (dms_reconnect_cycle_failures >= DMS_RECONNECT_MAX_CYCLE_FAILURES) {
k_sleep(K_SECONDS(2));
sys_reboot(SYS_REBOOT_COLD);
} else {
k_work_reschedule_for_queue(&dms_reconnect_work_q, &dms_reconnect_work,
K_SECONDS(DMS_RECONNECT_CYCLE_BACKOFF_SEC));
}
}
atomic_set(&dms_reconnect_in_progress, 0);
}
static void dms_reconnect_workq_ensure_started(void)
{
if (atomic_set(&dms_reconnect_workq_started, 1) == 0) {
k_work_queue_start(&dms_reconnect_work_q,
dms_reconnect_stack_area,
K_THREAD_STACK_SIZEOF(dms_reconnect_stack_area),
DMS_RECONNECT_PRIORITY, NULL);
k_thread_name_set(&dms_reconnect_work_q.thread, "dms_reconnect");
}
}
void dms_trigger_reconnect(void)
{
dms_reconnect_workq_ensure_started();
k_work_reschedule_for_queue(&dms_reconnect_work_q, &dms_reconnect_work, K_NO_WAIT);
}
#if defined(CONFIG_AWS_IOT_STATIC_IPV4)
static int broker_init(void)
{
struct sockaddr_in *broker4 =((struct sockaddr_in *)&broker);
inet_pton(AF_INET, CONFIG_AWS_IOT_STATIC_IPV4_ADDR, &broker4->sin_addr);
broker4->sin_family = AF_INET;
broker4->sin_port = htons(CONFIG_AWS_IOT_PORT);
return 0;
}
#else
static int broker_init(void)
{
int err;
struct addrinfo *result;
struct addrinfo *addr;
struct addrinfo hints = {.ai_family = AWS_AF_FAMILY,.ai_socktype = SOCK_STREAM};
err = getaddrinfo(aws_host_name_buf, NULL, &hints, &result);
if (err) return -ECHILD;
addr = result;
while (addr != NULL)
{
if ((addr->ai_addrlen == sizeof(struct sockaddr_in)) &&
(AWS_AF_FAMILY == AF_INET))
{
struct sockaddr_in *broker4 =((struct sockaddr_in *)&broker);
broker4->sin_addr.s_addr = ((struct sockaddr_in *)addr->ai_addr)->sin_addr.s_addr;
broker4->sin_family = AF_INET;
broker4->sin_port = htons(CONFIG_AWS_IOT_PORT);
break;
}
addr = addr->ai_next;
}
freeaddrinfo(result);
return err;
}
#endif
static int client_broker_init(struct mqtt_client *const client)
{
int err;
mqtt_client_init(client);
err = broker_init();
if (err) return err;
client->broker = &broker;
client->evt_cb = mqtt_evt_handler;
client->client_id.utf8 = client_id_buf;
client->client_id.size = strlen(client_id_buf);
client->password = NULL;
client->user_name = NULL;
client->protocol_version = MQTT_VERSION_3_1_1;
client->rx_buf = rx_buffer;
client->rx_buf_size = sizeof(rx_buffer);
client->tx_buf = tx_buffer;
client->tx_buf_size = sizeof(tx_buffer);
client->transport.type = MQTT_TRANSPORT_SECURE;
static sec_tag_t sec_tag_list[] = {CONFIG_AWS_IOT_SEC_TAG};
struct mqtt_sec_config *tls_cfg = &(client->transport).tls.config;
tls_cfg->peer_verify = 2;
tls_cfg->cipher_count = 0;
tls_cfg->cipher_list = NULL;
tls_cfg->sec_tag_count = ARRAY_SIZE(sec_tag_list);
tls_cfg->sec_tag_list = sec_tag_list;
tls_cfg->hostname = aws_host_name_buf;
tls_cfg->session_cache = TLS_SESSION_CACHE_DISABLED;
#if defined(CONFIG_AWS_IOT_PROVISION_CERTIFICATES)
err = certificates_provision_cattle();
if (err) return err;
#endif
return err;
}
static int connect_client(struct aws_iot_config *const config)
{
int err;
err = client_broker_init(&client);
if (err) return err;
err = mqtt_connect(&client);
if (err) return connect_error_translate(err);
if (config != NULL)
{
config->socket = client.transport.tls.sock;
}
return 0;
}
static int connection_poll_start(void)
{
if (atomic_get(&connection_poll_active)) return -EINPROGRESS;
atomic_set(&disconnect_requested, 0);
k_sem_give(&connection_poll_sem);
return 0;
}
int aws_iot_ping(void)
{
if (client.unacked_ping) return -ECONNRESET;
return mqtt_ping(&client);
}
int aws_iot_keepalive_time_left(void)
{
return mqtt_keepalive_time_left(&client);
}
int aws_iot_input(void)
{
return mqtt_input(&client);
}
int aws_iot_send(const struct aws_iot_data *const tx_data)
{
if (atomic_get(&fota_active)) {
LOG_WRN("FOTA in progress, dropping publish request.");
log_to_flash("FOTA in progress, dropping publish request.");
return -EBUSY;
}
struct aws_iot_data tx_data_pub = {
.ptr = tx_data->ptr,
.len = tx_data->len,
.qos = tx_data->qos,
.message_id = tx_data->message_id,
.retain_flag = tx_data->retain_flag,
.dup_flag = tx_data->dup_flag,
.topic.type = tx_data->topic.type,
.topic.str = tx_data->topic.str,
.topic.len = tx_data->topic.len};
struct mqtt_publish_param param;
param.message.topic.qos = tx_data_pub.qos;
if (tx_data_pub.topic.str != NULL) {
param.message.topic.topic.utf8 = tx_data_pub.topic.str;
param.message.topic.topic.size = tx_data_pub.topic.len;
} else {
param.message.topic.topic.utf8 = CONFIG_AWS_TOPIC_NAME;
param.message.topic.topic.size = strlen(CONFIG_AWS_TOPIC_NAME);
}
param.message.payload.data = tx_data_pub.ptr;
param.message.payload.len = tx_data_pub.len;
param.dup_flag = tx_data_pub.dup_flag;
param.retain_flag = tx_data_pub.retain_flag;
param.message_id = (tx_data_pub.message_id == 0) ? k_cycle_get_32() : tx_data_pub.message_id;
return mqtt_publish(&client, ¶m);
}
int aws_iot_disconnect(void)
{
atomic_set(&disconnect_requested, 1);
return mqtt_disconnect(&client);
}
int aws_iot_connect(struct aws_iot_config *const config)
{
int err;
if (IS_ENABLED(CONFIG_AWS_IOT_CONNECTION_POLL_THREAD))
{
err = connection_poll_start();
if (err) return err;
}
else
{
atomic_set(&disconnect_requested, 0);
err = connect_client(config);
if (err) return err;
atomic_set(&aws_iot_disconnected, 0);
}
return 0;
}
int aws_iot_init(const struct aws_iot_config *const config, aws_iot_evt_handler_t event_handler)
{
int err;
snprintf(client_id_buf_dms, sizeof(client_id_buf_dms), "dms-%s", MY_DEVICE_IMEI);
dms_handler_init();
err = aws_iot_set_host_name(cattle_host_name_buf, strlen(cattle_host_name_buf));
if (err) return err;
#if defined(CONFIG_AWS_IOT_BROKER_HOST_NAME_APP)
err = aws_iot_set_host_name(config->host_name, config->host_name_len);
#else
err = aws_iot_set_host_name(CONFIG_AWS_IOT_BROKER_HOST_NAME, strlen(CONFIG_AWS_IOT_BROKER_HOST_NAME));
#endif
if (err) return err;
module_evt_handler = event_handler;
return err;
}
void fetch_time_from_sntp(int *epochArr)
{
struct sntp_time sntp_time;
struct sockaddr_in addr;
int ret;
ret = net_ipaddr_parse(SNTP_SERVER, strlen(SNTP_SERVER),(struct sockaddr *)&addr);
if (ret < 0) return;
addr.sin_family = AF_INET;
addr.sin_port = htons(123);
ret = sntp_simple(SNTP_SERVER, 5000, &sntp_time);
if (ret == 0) {
uint32_t ts_low = (uint32_t)(sntp_time.seconds & 0xFFFFFFFF);
snprintf(current_ts, sizeof(current_ts), "%u", ts_low);
}
if (ret < 0) return;
time_t fetchedTime = (time_t)sntp_time.seconds;
dms_update_sntp_offset((int64_t)fetchedTime);
uint64_t maxTime = 1000000000;
for (int epochIdx = 0; epochIdx < 10; epochIdx++) {
epochArr[epochIdx] = (int)((fetchedTime / maxTime) % 10);
maxTime /= 10;
}
}
int aws_iot_subscription_topics_add(const struct aws_iot_topic_data *const topic_list, size_t list_count)
{
if ((topic_list == NULL) || (list_count == 0)) return -EINVAL;
if ((list_count + app_topic_data.list_count) > CONFIG_AWS_IOT_APP_SUBSCRIPTION_LIST_COUNT) return -ENOMEM;
for (size_t i = 0; i < list_count; i++) {
app_topic_data.list[app_topic_data.list_count].topic.utf8 = topic_list[i].str;
app_topic_data.list[app_topic_data.list_count].topic.size = topic_list[i].len;
app_topic_data.list[app_topic_data.list_count].qos = CONFIG_AWS_QOS;
app_topic_data.list_count++;
}
return 0;
}
#if defined(CONFIG_AWS_IOT_CONNECTION_POLL_THREAD)
void aws_iot_cloud_poll(void)
{
int err;
struct pollfd fds[2];
struct aws_iot_evt aws_iot_evt = {
.type = AWS_IOT_EVT_DISCONNECTED,
.data = {.err = AWS_IOT_DISCONNECT_MISC}
};
start:
k_sem_take(&connection_poll_sem, K_FOREVER);
while (k_sem_take(&connection_poll_sem, K_NO_WAIT) == 0) { }
if (atomic_get(&fota_active)) {
k_sleep(K_SECONDS(5));
goto start;
}
extern bool is_network_connected;
{
int gsm_wait = 0;
while (!is_network_connected && gsm_wait < 90) {
k_sleep(K_SECONDS(1));
gsm_wait++;
}
if (!is_network_connected) {
k_sleep(K_SECONDS(15));
k_sem_give(&connection_poll_sem);
goto start;
}
}
device_register_with_dms();
extern int epochTime[10];
if (!sntp_synced) {
for (int i = 0; i < 5 && epochTime[0] == 0; i++) {
fetch_time_from_sntp(epochTime);
if (epochTime[0] != 0) {
sntp_synced = true;
break;
}
k_sleep(K_SECONDS(5));
}
if (!sntp_synced) {
epochTime[0] = 1;
sntp_synced = true;
}
}
atomic_set(&connection_poll_active, 1);
err = connect_client_dms();
k_sleep(K_MSEC(2000));
err = connect_client(NULL);
if (err) goto reset;
atomic_set(&aws_iot_disconnected, 0);
while (true) {
if (atomic_get(&fota_active)) {
k_sleep(K_SECONDS(2));
continue;
}
int cattle_sock = client.transport.tls.sock;
if (cattle_sock < 0) break;
fds[0].fd = cattle_sock;
fds[0].events = POLLIN;
fds[1].fd = client_dms.transport.tls.sock;
fds[1].events = POLLIN;
int nfds = (client_dms.transport.tls.sock >= 0) ? 2 : 1;
int timeout = aws_iot_keepalive_time_left();
err = poll(fds, nfds, timeout);
if (err == 0) {
err = aws_iot_ping();
if (err) break;
if (nfds == 2) mqtt_ping(&client_dms);
continue;
}
if (err < 0) break;
if (fds[0].revents & POLLIN) {
err = aws_iot_input();
if (err || atomic_get(&aws_iot_disconnected) == 1) break;
}
if (nfds == 2 && (fds[1].revents & POLLIN)) {
err = mqtt_input(&client_dms);
if (err) {
k_mutex_lock(&dms_client_mutex, K_FOREVER);
mqtt_disconnect(&client_dms);
client_dms.transport.tls.sock = -1;
k_mutex_unlock(&dms_client_mutex);
}
}
if (fds[0].revents & (POLLERR | POLLHUP | POLLNVAL)) break;
if (nfds == 2 && (fds[1].revents & (POLLERR | POLLHUP | POLLNVAL))) {
k_mutex_lock(&dms_client_mutex, K_FOREVER);
mqtt_disconnect(&client_dms);
client_dms.transport.tls.sock = -1;
k_mutex_unlock(&dms_client_mutex);
}
}
if (atomic_get(&aws_iot_disconnected) == 0) {
aws_iot_notify_event(&aws_iot_evt);
aws_iot_disconnect();
}
if (atomic_get(&fota_active)) {
while (atomic_get(&fota_active)) { k_sleep(K_SECONDS(2)); }
atomic_set(&connection_poll_active, 0);
k_sem_take(&connection_poll_sem, K_NO_WAIT);
k_sleep(K_SECONDS(3));
k_sem_give(&connection_poll_sem);
goto start;
}
reset:
{
/* Same stale-fd fix as the FOTA teardown above: check the
* CURRENT socket value after mqtt_disconnect(), not one
* captured before it. */
k_mutex_lock(&dms_client_mutex, K_FOREVER);
if (client_dms.transport.tls.sock >= 0) {
mqtt_disconnect(&client_dms);
k_sleep(K_MSEC(100));
if (client_dms.transport.tls.sock >= 0) {
close(client_dms.transport.tls.sock);
}
}
client_dms.transport.tls.sock = -1;
k_mutex_unlock(&dms_client_mutex);
if (client.transport.tls.sock >= 0) {
mqtt_disconnect(&client);
k_sleep(K_MSEC(100));
if (client.transport.tls.sock >= 0) {
close(client.transport.tls.sock);
}
}
client.transport.tls.sock = -1;
}
atomic_set(&connection_poll_active, 0);
k_sem_take(&connection_poll_sem, K_NO_WAIT);
while (atomic_get(&fota_active)) { k_sleep(K_SECONDS(2)); }
k_sleep(K_SECONDS(20));
k_sem_give(&connection_poll_sem);
goto start;
}
K_THREAD_DEFINE(aws_connection_poll_thread, CONFIG_AWS_IOT_POLL_THREAD_STACK_SIZE, aws_iot_cloud_poll, NULL, NULL, NULL, K_PRIO_PREEMPT(0), 0, 0);
#endif/*
* node_fota_download.c — v4
*
* ADDED: a real retry loop around download_client_init()+get(), because
* connecting to S3 right after tearing down both MQTT TLS sockets was
* consistently failing with errno 113 (EHOSTUNREACH) / 118 on the first
* attempt — the modem/PPP link needs a moment to settle after two TLS
* sockets close, and a single extra k_sleep() wasn't reliably enough.
* Now it retries the CONNECT step itself, with increasing backoff,
* before giving up — this is the real "proper socket management" fix,
* not just a longer fixed delay.
*
* Everything else unchanged from v3: only ONE DMS metadata endpoint,
* differentiated by version string in iot_connections.c; this file
* just takes host/path and streams straight to flash, no dfu_target,
* no content checks of any kind.
*/
#include <zephyr/kernel.h>
#include <zephyr/logging/log.h>
#include <zephyr/drivers/flash.h>
#include <zephyr/device.h>
#include <zephyr/net/socket.h>
#include <zephyr/net/tls_credentials.h>
#include <net/download_client.h>
#include <mbedtls/sha256.h>
#include <string.h>
#include <errno.h>
#include <stdio.h>
#include "node_fota_download.h"
LOG_MODULE_REGISTER(node_fota_dl, LOG_LEVEL_INF);
#include "aws_iot/aws-certs.h"
#define SEC_TAG_S3 12
#define SPI_FLASH_SECTOR_SIZE 4096
/* ── Retry policy for the initial connect step ─────────────────────
* errno 113 (EHOSTUNREACH) / 118 seen right after MQTT teardown are
* transient — the PPP route/socket table on the modem hasn't fully
* settled yet. Retry the connect itself, not the whole DMS metadata
* round-trip (that's a much heavier, separate retry already handled
* by fota_retry_work in iot_connections.c). */
#define CONNECT_MAX_ATTEMPTS 5
static const int connect_backoff_sec[CONNECT_MAX_ATTEMPTS] = { 10, 15, 20, 25, 30 };
/* ── Dedicated DNS pre-check ─────────────────────────────────────
* download_client resolves the hostname internally too, but its
* failure just shows up as an opaque -118/EIO with no visibility.
* Resolving explicitly first, and logging the result, tells us
* definitively whether this is a DNS/network problem (this fails)
* or something else (this succeeds but download still fails). */
#define DNS_MAX_ATTEMPTS 4
static const int dns_backoff_sec[DNS_MAX_ATTEMPTS] = { 5, 8, 12, 15 };
static bool resolve_host_with_retry(const char *host, char *ip_out, size_t ip_out_len)
{
for (int attempt = 1; attempt <= DNS_MAX_ATTEMPTS; attempt++) {
struct addrinfo *res = NULL;
struct addrinfo hints = {
.ai_family = AF_INET,
.ai_socktype = SOCK_STREAM,
};
LOG_INF(">>> DNS attempt %d/%d for %s", attempt, DNS_MAX_ATTEMPTS, host);
int err = getaddrinfo(host, NULL, &hints, &res);
if (err == 0 && res != NULL) {
char ip_str[INET_ADDRSTRLEN] = {0};
struct sockaddr_in *sa = (struct sockaddr_in *)res->ai_addr;
inet_ntop(AF_INET, &sa->sin_addr, ip_str, sizeof(ip_str));
LOG_INF(">>> DNS resolved %s -> %s", host, ip_str);
if (ip_out) {
strncpy(ip_out, ip_str, ip_out_len - 1);
ip_out[ip_out_len - 1] = '\0';
}
freeaddrinfo(res);
return true;
}
LOG_ERR(">>> DNS resolution failed for %s (err=%d)", host, err);
if (attempt < DNS_MAX_ATTEMPTS) {
int backoff = dns_backoff_sec[attempt - 1];
LOG_WRN(">>> DNS retry in %ds...", backoff);
k_sleep(K_SECONDS(backoff));
}
}
return false;
}
extern atomic_t fota_active;
static struct download_client dlc;
static struct download_client_cfg dlc_cfg;
static struct {
char node_mac[32];
char host[128];
char path[3072];
size_t bytes_written;
mbedtls_sha256_context sha_ctx;
bool had_error;
node_fota_download_done_cb_t done_cb;
} ctx;
/* ── Diagnostic: is the PPP interface actually usable for a NEW TCP
* connection right now, or is this specific to the S3 destination?
* Tries a plain (non-TLS) connect to a host we already know works
* (the DMS server, which just had a live MQTT session moments ago).
* If THIS also fails with the same errno, the problem is the PPP
* interface/route state after MQTT teardown, not S3/DNS/the URL.
* If THIS succeeds but S3 still fails, the problem is specific to
* the S3 connection itself (port 443, TLS, or that specific route). */
static void diagnostic_control_connect(void)
{
struct sockaddr_in server = {0};
int sock = socket(AF_INET, SOCK_STREAM, IPPROTO_TCP);
if (sock < 0) {
LOG_ERR(">>> DIAG: control socket() failed, errno %d", errno);
return;
}
struct timeval tv = { .tv_sec = 5, .tv_usec = 0 };
setsockopt(sock, SOL_SOCKET, SO_SNDTIMEO, &tv, sizeof(tv));
server.sin_family = AF_INET;
server.sin_port = htons(3600);
net_addr_pton(AF_INET, "65.1.96.216", &server.sin_addr);
LOG_INF(">>> DIAG: control connect to known-good DMS server (65.1.96.216:3600)...");
int err = connect(sock, (struct sockaddr *)&server, sizeof(server));
if (err < 0) {
LOG_ERR(">>> DIAG: control connect FAILED, errno %d — PPP interface/route "
"itself is unusable right now, unrelated to S3/DNS/the URL", errno);
} else {
LOG_INF(">>> DIAG: control connect SUCCEEDED — PPP interface is fine, "
"problem is specific to the S3 connection");
}
close(sock);
}
/* ── Diagnostic #2: is a NEW TLS session specifically the problem
* (right after the two MQTT TLS sessions tore down), independent of
* download_client's own internal logic? Opens a raw TLS socket
* directly to the S3 IP:443 using the same SEC_TAG_S3 credentials,
* bypassing download_client entirely. If THIS also fails the same
* way, the problem is TLS session/context handling in general after
* MQTT teardown — not anything specific to download_client's own
* setup sequence. If THIS succeeds, the problem is something
* download_client itself does differently. */
static void diagnostic_tls_control_connect(const char *host_ip)
{
struct sockaddr_in server = {0};
int sock = socket(AF_INET, SOCK_STREAM, IPPROTO_TLS_1_2);
if (sock < 0) {
LOG_ERR(">>> DIAG-TLS: socket() failed, errno %d", errno);
return;
}
sec_tag_t sec_tag_list[] = { SEC_TAG_S3 };
int err = setsockopt(sock, SOL_TLS, TLS_SEC_TAG_LIST,
sec_tag_list, sizeof(sec_tag_list));
if (err < 0) {
LOG_ERR(">>> DIAG-TLS: setsockopt(TLS_SEC_TAG_LIST) failed, errno %d", errno);
close(sock);
return;
}
err = setsockopt(sock, SOL_TLS, TLS_HOSTNAME,
ctx.host, strlen(ctx.host) + 1);
if (err < 0) {
LOG_WRN(">>> DIAG-TLS: setsockopt(TLS_HOSTNAME) failed, errno %d "
"(continuing anyway)", errno);
}
struct timeval tv = { .tv_sec = 10, .tv_usec = 0 };
setsockopt(sock, SOL_SOCKET, SO_SNDTIMEO, &tv, sizeof(tv));
server.sin_family = AF_INET;
server.sin_port = htons(443);
net_addr_pton(AF_INET, host_ip, &server.sin_addr);
LOG_INF(">>> DIAG-TLS: raw TLS connect to %s:443 (bypassing download_client)...", host_ip);
err = connect(sock, (struct sockaddr *)&server, sizeof(server));
if (err < 0) {
LOG_ERR(">>> DIAG-TLS: raw TLS connect FAILED, errno %d — the problem is "
"TLS session handling itself right after MQTT teardown, not "
"download_client's specific setup", errno);
} else {
LOG_INF(">>> DIAG-TLS: raw TLS connect SUCCEEDED — TLS itself is fine, "
"something in download_client's own sequence is different");
}
close(sock);
}
static atomic_t dl_in_progress = ATOMIC_INIT(0);
static int flash_write_streaming(const struct device *flash_dev,
off_t base_addr, size_t offset,
const uint8_t *data, size_t len)
{
off_t abs_addr = base_addr + offset;
if (abs_addr % SPI_FLASH_SECTOR_SIZE == 0) {
int rc = flash_erase(flash_dev, abs_addr, SPI_FLASH_SECTOR_SIZE);
if (rc != 0) { LOG_ERR("flash_erase failed at 0x%lx: %d", (long)abs_addr, rc); return rc; }
}
int rc = flash_write(flash_dev, abs_addr, data, len);
if (rc != 0) LOG_ERR("flash_write failed at 0x%lx: %d", (long)abs_addr, rc);
return rc;
}
static int dlc_callback(const struct download_client_evt *event)
{
const struct device *flash_dev = DEVICE_DT_GET(DT_ALIAS(spi_flash0));
switch (event->id) {
case DOWNLOAD_CLIENT_EVT_FRAGMENT: {
size_t len = event->fragment.len;
if (ctx.bytes_written + len > NODE_FOTA_FLASH_MAX_SIZE) {
LOG_ERR(">>> Node image exceeds reserved flash region (%u bytes max)",
(unsigned)NODE_FOTA_FLASH_MAX_SIZE);
ctx.had_error = true;
return -ENOSPC;
}
int rc = flash_write_streaming(flash_dev, NODE_FOTA_FLASH_ADDR,
ctx.bytes_written,
event->fragment.buf, len);
if (rc != 0) { ctx.had_error = true; return rc; }
mbedtls_sha256_update(&ctx.sha_ctx, event->fragment.buf, len);
ctx.bytes_written += len;
if ((ctx.bytes_written % 8192) < len) {
LOG_INF(">>> Node firmware download progress: %u bytes written to flash",
(unsigned)ctx.bytes_written);
}
return 0;
}
case DOWNLOAD_CLIENT_EVT_DONE:
LOG_INF(">>> Node firmware download COMPLETE: %u bytes total",
(unsigned)ctx.bytes_written);
return 0;
case DOWNLOAD_CLIENT_EVT_ERROR:
LOG_ERR(">>> Node firmware download_client ERROR: %d", event->error);
ctx.had_error = true;
return 0;
default:
return 0;
}
}
#define DL_THREAD_STACK_SIZE 8192
static K_THREAD_STACK_DEFINE(dl_thread_stack, DL_THREAD_STACK_SIZE);
static struct k_thread dl_thread_data;
/* Returns 0 on a connect that at least STARTED (fragments may still
* fail later, handled by the caller's overall bytes_written check).
* Returns non-zero only if download_client_init/get itself rejected
* the attempt (the errno-113/118 case we're retrying around). */
static int try_connect_and_download_once(const char *resolved_ip)
{
int err;
err = download_client_init(&dlc, dlc_callback);
if (err) {
LOG_ERR(">>> download_client_init failed: %d", err);
download_client_disconnect(&dlc);
return err;
}
LOG_INF(">>> Connecting to %s (resolved: %s) ...", ctx.host, resolved_ip);
err = download_client_get(&dlc, resolved_ip, &dlc_cfg, ctx.path, 0);
if (err) {
LOG_ERR(">>> download_client_get failed: %d", err);
download_client_disconnect(&dlc);
return err;
}
return 0;
}
static void dl_thread_fn(void *a, void *b, void *c)
{
ARG_UNUSED(a); ARG_UNUSED(b); ARG_UNUSED(c);
uint8_t sha_out[32];
bool success = false;
int err = -1;
LOG_INF("========================================");
LOG_INF(" NODE FIRMWARE DOWNLOAD STARTING");
LOG_INF(" Host : %s", ctx.host);
LOG_INF(" Flash: 0x%lx (max %u bytes)",
(long)NODE_FOTA_FLASH_ADDR, (unsigned)NODE_FOTA_FLASH_MAX_SIZE);
LOG_INF("========================================");
log_to_flash("Node firmware: download thread started");
ctx.bytes_written = 0;
ctx.had_error = false;
mbedtls_sha256_init(&ctx.sha_ctx);
mbedtls_sha256_starts(&ctx.sha_ctx, 0);
tls_credential_delete(SEC_TAG_S3, TLS_CREDENTIAL_CA_CERTIFICATE);
tls_credential_add(SEC_TAG_S3, TLS_CREDENTIAL_CA_CERTIFICATE,
amazon_root_ca_1, amazon_root_ca1_len);
dlc_cfg = (struct download_client_cfg){
.sec_tag_list = (sec_tag_t[]){ SEC_TAG_S3 },
.sec_tag_count = 1,
.set_tls_hostname = ctx.host, /* SNI hostname for TLS handshake */
};
char resolved_ip[INET_ADDRSTRLEN] = {0};
if (!resolve_host_with_retry(ctx.host, resolved_ip, sizeof(resolved_ip))) {
LOG_ERR(">>> Giving up: could not resolve %s after %d attempts. "
"This is a network/DNS problem, not a URL parsing problem "
"(the host string itself is correct — see log above).",
ctx.host, DNS_MAX_ATTEMPTS);
log_to_flash("Node firmware: DNS resolution failed after retries");
goto done;
}
/* ── Run diagnostics NOW, while we have the resolved IP ── */
diagnostic_control_connect(); /* plain TCP to DMS — is PPP usable at all? */
diagnostic_tls_control_connect(resolved_ip); /* raw TLS to S3 — is TLS the problem? */
LOG_INF(">>> Waiting for stack to settle after MQTT teardown...");
k_sleep(K_SECONDS(5));
for (int attempt = 1; attempt <= CONNECT_MAX_ATTEMPTS; attempt++) {
ctx.bytes_written = 0;
ctx.had_error = false;
LOG_INF(">>> Connect attempt %d/%d", attempt, CONNECT_MAX_ATTEMPTS);
err = try_connect_and_download_once(resolved_ip);
if (err == 0) {
/* Connect accepted — now wait for the transfer to actually
* progress or fail via the callback. */
int idle_ms = 0;
size_t last_seen = 0;
while (idle_ms < 15000 && !ctx.had_error) {
k_sleep(K_MSEC(500));
if (ctx.bytes_written == last_seen) {
idle_ms += 500;
} else {
idle_ms = 0;
last_seen = ctx.bytes_written;
}
}
download_client_disconnect(&dlc);
if (!ctx.had_error && ctx.bytes_written > 0) {
break; /* success, drop out of retry loop */
}
LOG_WRN(">>> Transfer failed/stalled after connect (bytes=%u), "
"will retry connect", (unsigned)ctx.bytes_written);
}
if (attempt < CONNECT_MAX_ATTEMPTS) {
int backoff = connect_backoff_sec[attempt - 1];
LOG_WRN(">>> Connect attempt %d failed, retrying in %ds...",
attempt, backoff);
log_to_flash("Node firmware: connect attempt failed, retrying");
k_sleep(K_SECONDS(backoff));
}
}
if (ctx.had_error || ctx.bytes_written == 0) {
LOG_ERR(">>> Node firmware download FAILED after %d attempts",
CONNECT_MAX_ATTEMPTS);
log_to_flash("Node firmware: download failed after all retries");
goto done;
}
mbedtls_sha256_finish(&ctx.sha_ctx, sha_out);
success = true;
LOG_INF(">>> Node firmware download SUCCESS: %u bytes, SHA-256 computed",
(unsigned)ctx.bytes_written);
log_to_flash("Node firmware: download succeeded");
done:
mbedtls_sha256_free(&ctx.sha_ctx);
if (ctx.done_cb) {
ctx.done_cb(success, ctx.node_mac, NODE_FOTA_FLASH_ADDR,
ctx.bytes_written, success ? sha_out : NULL);
}
atomic_set(&dl_in_progress, 0);
}
static void unescape_buf(char *str) {
char *src = str;
char *dst = str;
while (*src) {
if (*src == '\\') {
src++; // Skip the backslash
if (*src == '/') {
*dst++ = '/';
src++;
} else if (*src == 'u') {
// Handle \u0026 (&) and general 4-digit hex unicode escapes
src++;
char hex[5] = {0};
if (strlen(src) >= 4) {
memcpy(hex, src, 4);
unsigned int code = (unsigned int)strtoul(hex, NULL, 16);
*dst++ = (char)code;
src += 4;
}
} else if (*src == 'n') {
*dst++ = '\n';
src++;
} else if (*src == 'r') {
*dst++ = '\r';
src++;
} else if (*src == 't') {
*dst++ = '\t';
src++;
} else {
*dst++ = *src++;
}
} else {
*dst++ = *src++;
}
}
*dst = '\0'; // Null-terminate
}
bool node_fota_download_in_progress(void) { return atomic_get(&dl_in_progress) != 0; }
int node_fota_download_from_host_path(const char *node_mac,
const char *host, const char *path,
node_fota_download_done_cb_t done_cb)
{
if (atomic_cas(&dl_in_progress, 0, 1) == false) {
LOG_WRN(">>> Node FOTA download already in progress, rejecting new request");
return -EBUSY;
}
strncpy(ctx.node_mac, node_mac, sizeof(ctx.node_mac) - 1);
ctx.node_mac[sizeof(ctx.node_mac) - 1] = '\0';
strncpy(ctx.host, host, sizeof(ctx.host) - 1);
ctx.host[sizeof(ctx.host) - 1] = '\0';
unescape_buf(ctx.host); // <--- UNESCAPE HOST
strncpy(ctx.path, path, sizeof(ctx.path) - 1);
ctx.path[sizeof(ctx.path) - 1] = '\0';
unescape_buf(ctx.path); // <--- UNESCAPE PATH
ctx.done_cb = done_cb;
k_thread_create(&dl_thread_data, dl_thread_stack,
K_THREAD_STACK_SIZEOF(dl_thread_stack),
dl_thread_fn, NULL, NULL, NULL,
8, 0, K_NO_WAIT);
return 0;
}