#include "local_data_lasting.h"
#include <sys/mman.h>
#include <pthread.h>
#include <errno.h>
#include <stdio.h>
#include <string.h>
#include <stddef.h>
#include "protocol.h"
#include "common.h"
#include "proto_forward.h"
#include "tag.h"
#include "define.h"
#include "collector-api.h"

#define SYS_PERSISTENCE_DB_TB "data_per"
#define NEW_TAG_SIZE 50
#define DEF_SHM_TAG_NUM 100

sqlite3 *create_presistence_table(const char *db_path);
int is_database_file_exists(const char *db_path);
static void str_replace_chr(char *line);
static int shared_memory_restore_from_db(data_presistence_info * presistence_ptr);
sqlite3 *open_existing_database(const char *db_path);
extern proto_forward_t *get_proto_forward_var(void);
static int calc_tag_capacity(size_t shm_size) {
    return (int)((shm_size - offsetof(shm_info_t, tags)) / sizeof(shm_tag_info));
}

void updata_data_persistence(void *par_tar) {
    tag_t *tag = (tag_t *)par_tar;
    if (tag->shm_tag_p == NULL)
    {
        ems_syslog(LOG_ERR, "tag[%s] shm_tag_p is null", tag->dev_tag_name);
        return;
    }
    proto_syslog(LOG_DEBUG, "[updata_data_persistence]: tag[%d]=%s, %ld, %f", tag->data_type, tag->dev_tag_name, tag->read_cache.to_int, tag->read_cache.to_float);
    if (tag->data_type == 0) {
        tag->shm_tag_p->tag_cache.to_int = tag->read_cache.to_int;
    }
    else {
        tag->shm_tag_p->tag_cache.to_float = tag->read_cache.to_float;
    }
}

sqlite3 *open_existing_database(const char *db_path) {
    sqlite3 *db = NULL;
    int rc;
    rc = sqlite3_open_v2(db_path, &db, SQLITE_OPEN_READWRITE, NULL);
    if (rc != SQLITE_OK) {
        proto_syslog(LOG_ERR, "Failed to open database: %s", sqlite3_errmsg(db));
        sqlite3_close(db);
        return NULL;
    }
    return db;
}

int expand_shm_info(data_presistence_info *presistence_ptr, size_t new_count)
{
    if (!presistence_ptr || !presistence_ptr->shm_info)
        return -1;

    shm_info_t *shm_ptr = presistence_ptr->shm_info;
    size_t old_shm_size = sizeof(shm_info_t) + shm_ptr->count * sizeof(shm_tag_info);
    size_t new_shm_size = sizeof(shm_info_t) + new_count * sizeof(shm_tag_info);
    int fd = shm_ptr->fd;
    shm_info_t *back_shm = calloc(1, old_shm_size);
    if (back_shm == NULL)
    {
        return -10;
    }
    //
    memcpy(back_shm, shm_ptr, old_shm_size);
    //
    int old_tags_index = shm_ptr->tags_index;

    if (munmap(shm_ptr, old_shm_size) == -1)
    {
        proto_syslog(LOG_ERR, "munmap failed: %s", strerror(errno));
        return -1;
    }
    if (ftruncate(fd, new_shm_size) == -1)
    {
        proto_syslog(LOG_ERR, "ftruncate failed: %s", strerror(errno));
        return -1;
    }
    shm_info_t *new_shm_ptr = mmap(NULL, new_shm_size, PROT_READ | PROT_WRITE, MAP_SHARED, fd, 0);
    if (new_shm_ptr == MAP_FAILED)
    {
        proto_syslog(LOG_ERR, "mmap failed after expand: %s", strerror(errno));
        return -1;
    }
    memset(new_shm_ptr, 0, new_shm_size);
    presistence_ptr->shm_info = new_shm_ptr;

    int tag_capacity = calc_tag_capacity(new_shm_size);
    new_shm_ptr->fd = fd;
    new_shm_ptr->count = tag_capacity;
    new_shm_ptr->tags_index = old_tags_index;
    for (int i = 0; i < new_shm_ptr->tags_index; i++)
    {
        new_shm_ptr->tags[i] = back_shm->tags[i];
        if (new_shm_ptr->tags[i].tag_cache_p != NULL)
        {
            ((tag_t*)new_shm_ptr->tags[i].tag_cache_p)->shm_tag_p = &new_shm_ptr->tags[i];
        }
    }
    free(back_shm);

    proto_syslog(LOG_WARNING, "Shared memory expanded, new count: %d, tags_index: %d",
        tag_capacity, new_shm_ptr->tags_index);

    return 0;
}

static int file_exists(const char *path) {
    struct stat st;
    return stat(path, &st) == 0;
}

static int shm_exists(const char *shm_path) {
    return file_exists(shm_path);
}

sqlite3 *create_presistence_table(const char *db_path)
{
    sqlite3 *db = NULL;
    char *zErrMsg = NULL;
    char sql[512];
    char name[64] = {0};

    local_strlcpy(name, SYS_PERSISTENCE_DB_TB, sizeof(name) - 1);
    str_replace_chr(name);

    int rc = sqlite3_open(db_path, &db);
    if (rc != SQLITE_OK)
    {
        proto_syslog(LOG_ERR, "dbopen failed：%s", sqlite3_errmsg(db));
        sqlite3_close(db);
        return NULL;
    }

    snprintf(sql, sizeof(sql),
             "CREATE TABLE IF NOT EXISTS %s ("
             "ID INTEGER PRIMARY KEY AUTOINCREMENT, "
             "status INTEGER NOT NULL DEFAULT 0, "
             "dev_tag_name VARCHAR(96) NOT NULL, "
             "tag_cache_int INTEGER, "
             "tag_cache_float REAL, "
             "datatype INTEGER);",
             name);

    proto_syslog(LOG_DEBUG, "create table：%s", name);
    rc = sqlite3_exec(db, sql, NULL, NULL, &zErrMsg);
    if (rc != SQLITE_OK)
    {
        proto_syslog(LOG_ERR, "craete table failed：%s，error：%s", name, zErrMsg);
        sqlite3_free(zErrMsg);
        sqlite3_close(db);
        return NULL;
    }

    snprintf(sql, sizeof(sql),
             "CREATE INDEX IF NOT EXISTS dev_tag_index ON %s (dev_tag_name);",
             name);
    rc = sqlite3_exec(db, sql, NULL, NULL, &zErrMsg);
    if (rc != SQLITE_OK)
    {
        proto_syslog(LOG_ERR, "create failed：%s，error：%s", name, zErrMsg);
        sqlite3_free(zErrMsg);
        sqlite3_close(db);
        return NULL;
    }

    return db;
}

shm_info_t *create_shm_memory(const char *shm_path, size_t shm_size, int is_new)
{
    int fd = -1;
    shm_info_t *shm_ptr = NULL;
    if (!shm_path || shm_size == 0) {
        proto_syslog(LOG_ERR, "Invalid shm_path or shm_size");
        return NULL;
    }

    fd = open(shm_path, is_new ? (O_RDWR | O_CREAT) : O_RDWR, 0666);
    if (fd < 0) {
        proto_syslog(LOG_ERR, "open shm failed: %s", strerror(errno));
        return NULL;
    }
    if (is_new) {
        if (ftruncate(fd, shm_size) == -1) {
            proto_syslog(LOG_ERR, "ftruncate failed: %s", strerror(errno));
            close(fd);
            return NULL;
        }
    }
    shm_ptr = (shm_info_t *)mmap(NULL, shm_size, PROT_READ | PROT_WRITE, MAP_SHARED, fd, 0);
    if (shm_ptr == MAP_FAILED) {
        proto_syslog(LOG_ERR, "mmap failed: %s", strerror(errno));
        close(fd);
        return NULL;
    }
    if (is_new)
    {
        memset(shm_ptr, 0, shm_size);
        int tag_capacity = calc_tag_capacity(shm_size);
        shm_ptr->count = tag_capacity;
    }
    else
    {
        int count = shm_ptr->count;
        proto_syslog(LOG_NOTICE, "create_shm_memory: count:%d, tags_index:%d", shm_ptr->count, shm_ptr->tags_index);
        if (munmap(shm_ptr, shm_size) == -1)
        {
            proto_syslog(LOG_ERR, "munmap failed: %s", strerror(errno));
            return NULL;
        }
        shm_size = sizeof(shm_info_t) + sizeof(shm_tag_info) * count;
        shm_ptr = (shm_info_t *)mmap(NULL, shm_size, PROT_READ | PROT_WRITE, MAP_SHARED, fd, 0);
        if (shm_ptr == MAP_FAILED)
        {
            proto_syslog(LOG_ERR, "mmap failed: %s", strerror(errno));
            close(fd);
            return NULL;
        }
    }
    shm_ptr->fd = fd;
    for (int i = 0; i < shm_ptr->count; i++)
    {
        shm_ptr->tags[i].tag_cache_p = NULL;
    }
    proto_syslog(LOG_ERR, "create_shm_memory: shm_path=%s, size=%zu, is_new=%d, count=%d", shm_path, shm_size, is_new, shm_ptr->count);

    return shm_ptr;
}

static sqlite3 *new_db(const char *db_path)
{
    int db_exist = file_exists(db_path);
    if (db_exist != 0)
    {
        return open_existing_database(db_path);
    }
    else
    {
        return create_presistence_table(db_path);
    }
}

data_presistence_info *new_data_presistence(const char *shm_name, int update_interval) {
    if (!shm_name || update_interval <= 0) {
        proto_syslog(LOG_ERR, "Invalid parameters: shm_name or update_interval");
        return NULL;
    }

    proto_syslog(LOG_INFO, "[STARTUP] Starting data persistence module, name=%s, interval=%d", shm_name, update_interval);

    data_presistence_info *info = (data_presistence_info *)malloc(sizeof(data_presistence_info));
    if (!info) {
        proto_syslog(LOG_ERR, "[STARTUP] data_presistence_info malloc failed");
        return NULL;
    }
    memset(info, 0, sizeof(data_presistence_info));
    info->db_update_interval = update_interval;

    char db_path[256];
    char shm_path[256];
    snprintf(db_path, sizeof(db_path), "%s/%s.db", SYS_PERSISTENCE_DIR, shm_name);
    snprintf(shm_path, sizeof(shm_path), "/dev/shm/%s", shm_name);

    info->init_flag = 0;
    size_t shm_size = sizeof(shm_info_t) + sizeof(shm_tag_info) * DEF_SHM_TAG_NUM;
    info->db_info = new_db(db_path);
    int shm_exist = shm_exists(shm_path);
    if (shm_exist != 0) // 共享内存存在
    {
        info->shm_info = create_shm_memory(shm_path, shm_size, 0);
    }
    else
    {
        info->shm_info = create_shm_memory(shm_path, shm_size, 1);
        shared_memory_restore_from_db(info);
    }
    shm_tag_info *p = NULL;
    for (int i = 0; i < info->shm_info->tags_index; i++)
    {
        p = &info->shm_info->tags[i];
        proto_syslog(LOG_NOTICE, "%p tag_name:%s, tag_cache_p:%p\n", p, p->dev_tag_name, p->tag_cache_p);
    }

    info->init_flag = 1;
    proto_syslog(LOG_DEBUG, "new_data_presistence init done, shm=%s, db=%s", shm_path, db_path);
    return info;
}

void get_last_data_from_data_presistence(data_presistence_info * presistence_ptr, void *param)
{
    tag_t *tag = (tag_t*)param;
    shm_info_t *shm_ptr = presistence_ptr->shm_info;
    if (!shm_ptr) {
        proto_syslog(LOG_ERR, "[TAG] Shared memory pointer is null in get_last_data_from_data_presistence");
        return;
    }

    int found = 0;

    for (int i = 0; i < shm_ptr->tags_index; i++)
    {
        if (shm_ptr->tags_index >= shm_ptr->count)
        {
            size_t new_count = shm_ptr->count + NEW_TAG_SIZE;
            if (expand_shm_info(presistence_ptr, new_count) != 0)
            {
                return;
            }
            shm_ptr = presistence_ptr->shm_info;
            if (shm_ptr->tags_index >= shm_ptr->count) {
                proto_syslog(LOG_ERR, "tags_index(%d) >= count(%d) after expand!", shm_ptr->tags_index, shm_ptr->count);
                return;
            }
        }

        if (strcmp(tag->dev_tag_name, shm_ptr->tags[i].dev_tag_name) == 0)
        {
            tag->shm_tag_p = &shm_ptr->tags[i];
            shm_ptr->tags[i].tag_cache_p = tag;
            shm_ptr->tags[i].status = 1; // 当前配置了这个点位
            if (tag->data_type == 0)
                tag->read_cache.to_int = shm_ptr->tags[i].tag_cache.to_int;
            else
                tag->read_cache.to_float = shm_ptr->tags[i].tag_cache.to_float;

            found = 1;
            proto_syslog(LOG_INFO, "[TAG] Found tag %s at index %d", tag->dev_tag_name, i);
            break;
        }
    }

    if (!found)
    {
        if (shm_ptr->tags_index >= shm_ptr->count)
        {
            size_t new_count = shm_ptr->count + NEW_TAG_SIZE;
            if (expand_shm_info(presistence_ptr, new_count) != 0)
            {
                return;
            }
            shm_ptr = presistence_ptr->shm_info;
            if (shm_ptr->tags_index >= shm_ptr->count) {
                proto_syslog(LOG_ERR, "tags_index(%d) >= count(%d) after expand!", shm_ptr->tags_index, shm_ptr->count);
                return;
            }
        }
        shm_tag_info *new_tag = &shm_ptr->tags[shm_ptr->tags_index];
        new_tag->status = 1;
        snprintf(new_tag->dev_tag_name, sizeof(new_tag->dev_tag_name), "%s", tag->dev_tag_name);
        new_tag->datatype = tag->data_type;
        if (tag->data_type == 0)
            new_tag->tag_cache.to_int = tag->read_cache.to_int;
        else
            new_tag->tag_cache.to_float = tag->read_cache.to_float;

        tag->shm_tag_p = new_tag;
        new_tag->tag_cache_p = tag;
        shm_ptr->tags_index++;
        proto_syslog(LOG_INFO, "[TAG] Added new tag %s at index %d, %p", tag->dev_tag_name, shm_ptr->tags_index-1, tag->shm_tag_p);
    }
}

int shared_memory_restore_from_db(data_presistence_info * presistence_ptr)
{
    shm_info_t *global_shm_ptr = presistence_ptr->shm_info;
    if ((global_shm_ptr == NULL) || (presistence_ptr->db_info == NULL))
    {
        proto_syslog(LOG_ERR, "[RESTORE] Invalid parameters (shm or db is NULL)");
        return -1;
    }

    proto_syslog(LOG_INFO, "[RESTORE] Begin restoring shared memory from DB");

    sqlite3 *db = presistence_ptr->db_info;
    const char *sql = "SELECT status, dev_tag_name, tag_cache_int, tag_cache_float, datatype FROM " SYS_PERSISTENCE_DB_TB ";";
    sqlite3_stmt *stmt = NULL;
    int rc;
    rc = sqlite3_prepare_v2(db, sql, -1, &stmt, NULL);
    if (rc != SQLITE_OK)
    {
        proto_syslog(LOG_ERR, "[RESTORE] Prepare failed: %s", sqlite3_errmsg(db));
        return -1;
    }
    global_shm_ptr->tags_index = 0;
    int restore_count = 0;

    while ((rc = sqlite3_step(stmt)) == SQLITE_ROW)
    {
        int status = sqlite3_column_int(stmt, 0);
        const char *dev_tag_name = (const char *)sqlite3_column_text(stmt, 1);
        long tag_cache_int = sqlite3_column_int64(stmt, 2);
        double tag_cache_float = sqlite3_column_double(stmt, 3);
        int datatype = sqlite3_column_int(stmt, 4);

        if (global_shm_ptr->tags_index >= global_shm_ptr->count)
        {
            size_t new_count = global_shm_ptr->count + NEW_TAG_SIZE;
            if (expand_shm_info(presistence_ptr, new_count) != 0) {
                sqlite3_finalize(stmt);
                return -1;
            }
            global_shm_ptr = presistence_ptr->shm_info;
            if (global_shm_ptr->tags_index >= global_shm_ptr->count) {
                proto_syslog(LOG_ERR, "tags_index(%d) >= count(%d) after expand!", global_shm_ptr->tags_index, global_shm_ptr->count);
                sqlite3_finalize(stmt);
                return -1;
            }
        }

        shm_tag_info *tag = &global_shm_ptr->tags[global_shm_ptr->tags_index];
        tag->status = status;
        strncpy(tag->dev_tag_name, dev_tag_name, sizeof(tag->dev_tag_name) - 1);
        tag->dev_tag_name[sizeof(tag->dev_tag_name) - 1] = '\0';
        tag->datatype = datatype;
        tag->tag_cache_p = NULL;

        if (datatype == 0) {
            tag->tag_cache.to_int = tag_cache_int;
        } else {
            tag->tag_cache.to_float = tag_cache_float;
        }

        global_shm_ptr->tags_index++;
        restore_count++;
        proto_syslog(LOG_INFO, "[RESTORE] Restored tag %s at index %d", dev_tag_name, global_shm_ptr->tags_index-1);
    }

    if (rc != SQLITE_DONE)
    {
        proto_syslog(LOG_ERR, "[RESTORE] Query error: %s", sqlite3_errmsg(db));
        sqlite3_finalize(stmt);
        return -1;
    }

    sqlite3_finalize(stmt);
    proto_syslog(LOG_INFO, "[RESTORE] Finished restore, total tags restored: %d", restore_count);
    return 0;
}

int updata_sql_database(data_presistence_info * presistence_ptr)
{
    sqlite3 *db = presistence_ptr->db_info;
    shm_info_t *global_shm_ptr = presistence_ptr->shm_info;
    if (!db || !global_shm_ptr)
    {
        proto_syslog(LOG_ERR, "Invalid parameters");
        return -1;
    }
    struct prepare_v2{
        sqlite3_stmt *stmt;
        const char *sql;
    };

    struct prepare_v2 stmt_sql[] = {
        {.stmt = NULL, .sql = "SELECT tag_cache_int, tag_cache_float, datatype FROM " SYS_PERSISTENCE_DB_TB " WHERE dev_tag_name = ?;"},
        {.stmt = NULL, .sql = "UPDATE "SYS_PERSISTENCE_DB_TB" SET status = ?, tag_cache_int = ?, tag_cache_float = ? WHERE dev_tag_name = ?;"},
        {.stmt = NULL, .sql = "INSERT INTO " SYS_PERSISTENCE_DB_TB " (status, dev_tag_name, tag_cache_int, tag_cache_float, datatype) VALUES (?, ?, ?, ?, ?);"},
    };

    for (int i = 0; i < sizeof(stmt_sql) / sizeof(stmt_sql[0]); i++)
    {
        if (sqlite3_prepare_v2(db, stmt_sql[i].sql, -1, &stmt_sql[i].stmt, NULL) != SQLITE_OK)
        {
            proto_syslog(LOG_ERR, "Prepare failed[%d]: %s", i, sqlite3_errmsg(db));
            return -2;
        }
    }
    int rc = 0;
    sqlite3_stmt *stmt = NULL;
    for (int i = 0; i < global_shm_ptr->tags_index; ++i)
    {
        stmt = stmt_sql[0].stmt;
        shm_tag_info *tag = &global_shm_ptr->tags[i];
        sqlite3_bind_text(stmt, 1, tag->dev_tag_name, -1, SQLITE_STATIC);
        rc = sqlite3_step(stmt);
        if (rc == SQLITE_ROW)
        {
            long db_int = sqlite3_column_int64(stmt, 0);
            double db_float = sqlite3_column_double(stmt, 1);

            int need_update = 0;
            if (tag->datatype == 0)
            {
                if (tag->tag_cache.to_int != db_int)
                {
                    need_update = 1;
                }
            }
            else if (tag->datatype == 1)
            {
                if (tag->tag_cache.to_float != db_float)
                {
                    need_update = 1;
                }
            }
            if (need_update)
            {
                sqlite3_reset(stmt);
                stmt = stmt_sql[1].stmt;
                sqlite3_bind_int(stmt, 1, tag->status);
                if (tag->datatype == 0)
                {
                    sqlite3_bind_int64(stmt, 2, tag->tag_cache.to_int);
                    sqlite3_bind_null(stmt, 3);
                }
                else
                {
                    sqlite3_bind_null(stmt, 2);
                    sqlite3_bind_double(stmt, 3, tag->tag_cache.to_float);
                }
                sqlite3_bind_text(stmt, 4, tag->dev_tag_name, -1, SQLITE_STATIC);

                if (sqlite3_step(stmt) != SQLITE_DONE)
                {
                    proto_syslog(LOG_ERR, "Update failed: %s", sqlite3_errmsg(db));
                }
                else
                {
                    proto_syslog(LOG_DEBUG, "Updated: %s", tag->dev_tag_name);
                }
            }
        }
        else if (rc == SQLITE_DONE)
        {
            stmt = stmt_sql[2].stmt;
            sqlite3_bind_int(stmt, 1, tag->status);
            sqlite3_bind_text(stmt, 2, tag->dev_tag_name, -1, SQLITE_STATIC);
            if (tag->datatype == 0)
            {
                sqlite3_bind_int64(stmt, 3, tag->tag_cache.to_int);
                sqlite3_bind_null(stmt, 4);
            }
            else
            {
                sqlite3_bind_null(stmt, 3);
                sqlite3_bind_double(stmt, 4, tag->tag_cache.to_float);
            }
            sqlite3_bind_int(stmt, 5, tag->datatype);

            if (sqlite3_step(stmt) != SQLITE_DONE)
            {
                proto_syslog(LOG_ERR, "Insert failed: %s", sqlite3_errmsg(db));
            }
            else
            {
                //proto_syslog(LOG_DEBUG, "Inserted: %s", tag->dev_tag_name);
            }
        }
        sqlite3_reset(stmt);
    }

    for (int i = 0; i < sizeof(stmt_sql) / sizeof(stmt_sql[0]); i++)
    {
        sqlite3_finalize(stmt_sql[i].stmt);
    }
    return 1;
}

void *data_presistence_thread_function(void* param)
{
    data_presistence_info *presistence_ptr = (data_presistence_info*)param;

    while (presistence_ptr->db_info == NULL || presistence_ptr->shm_info == NULL || presistence_ptr->init_flag == 0)
    {
        sleep(1);
    }

    while (1)
    {
        updata_sql_database(presistence_ptr);
        ems_syslog(LOG_NOTICE, "updata_sql_database data_presistence_db");
        sleep(presistence_ptr->db_update_interval);
    }
}

int data_presistence_start(data_presistence_info * presistence_ptr, const char *thread_name)
{
    pthread_t data_presistence_thread_id = PTHREAD_NULL;
    int ret = pthread_create(&data_presistence_thread_id, NULL, data_presistence_thread_function, presistence_ptr);
    if (ret == 0)
    {
        pthread_setname_np(data_presistence_thread_id, thread_name);
    }
    return ret;
}

static void str_replace_chr(char *line)
{
    for (int i = 0; line[i] != '\0'; i++)
    {
        if (!((line[i] >= '0' && line[i] <= '9') || (line[i] >= 'a' && line[i] <= 'z') || (line[i] >= 'A' && line[i] <= 'Z') || line[i] == '\0' || line[i] == '_'))
        {
            line[i] = '_';
        }
    }
}