omnix-sopiga/services/omnix-broadcast/migrations/001_create_collection_broadcasts.up.sql
2026-08-07 12:21:42 +07:00

259 lines
9.9 KiB
PL/PgSQL

CREATE SCHEMA IF NOT EXISTS collection_broadcasts;
-- ============================================================================
-- 1. CONFIGURATION TABLES (Reference only)
-- ============================================================================
CREATE TABLE collection_broadcasts.sopiga_template_config (
id BIGSERIAL PRIMARY KEY,
template_name VARCHAR(100) NOT NULL,
sopiga_template_id INT NOT NULL UNIQUE,
template_type VARCHAR(50),
channel VARCHAR(50),
description TEXT,
active BOOLEAN DEFAULT TRUE,
created_at TIMESTAMP DEFAULT NOW(),
updated_at TIMESTAMP DEFAULT NOW()
);
CREATE INDEX idx_sopiga_template_active ON collection_broadcasts.sopiga_template_config(active);
CREATE TABLE collection_broadcasts.template_variable_mapping (
id BIGSERIAL PRIMARY KEY,
sopiga_template_id INT NOT NULL REFERENCES collection_broadcasts.sopiga_template_config(sopiga_template_id) ON DELETE CASCADE,
variable_order INT NOT NULL,
sopiga_variable_name VARCHAR(100) NOT NULL,
variable_type VARCHAR(50),
db_field_source VARCHAR(100) NOT NULL,
is_required BOOLEAN DEFAULT TRUE,
example_value VARCHAR(500),
description TEXT,
created_at TIMESTAMP DEFAULT NOW(),
updated_at TIMESTAMP DEFAULT NOW(),
CONSTRAINT unique_template_variable UNIQUE (sopiga_template_id, sopiga_variable_name),
CONSTRAINT unique_variable_order UNIQUE (sopiga_template_id, variable_order)
);
CREATE INDEX idx_template_variable_mapping_template_id ON collection_broadcasts.template_variable_mapping(sopiga_template_id);
CREATE INDEX idx_template_variable_mapping_order ON collection_broadcasts.template_variable_mapping(sopiga_template_id, variable_order);
CREATE TABLE collection_broadcasts.sopiga_collar_config (
id BIGSERIAL PRIMARY KEY,
broadcast_name VARCHAR(100) NOT NULL,
sopiga_collar_id INT NOT NULL UNIQUE,
sopiga_template_id INT NOT NULL REFERENCES collection_broadcasts.sopiga_template_config(sopiga_template_id),
status VARCHAR(50),
description TEXT,
created_at TIMESTAMP DEFAULT NOW(),
updated_at TIMESTAMP DEFAULT NOW()
);
CREATE INDEX idx_sopiga_collar_status ON collection_broadcasts.sopiga_collar_config(status);
-- ============================================================================
-- 2. BROADCAST STAGING (CORE) - DENORMALIZED
-- ============================================================================
CREATE TABLE collection_broadcasts.broadcast_staging (
id BIGSERIAL PRIMARY KEY,
sopiga_collar_id INT NOT NULL REFERENCES collection_broadcasts.sopiga_collar_config(sopiga_collar_id),
sopiga_template_id INT NOT NULL REFERENCES collection_broadcasts.sopiga_template_config(sopiga_template_id),
-- Single source of truth: all fields needed to build the message live here
message_payload JSONB NOT NULL,
status VARCHAR(50) DEFAULT 'pending' NOT NULL,
error_message TEXT,
error_count INT DEFAULT 0,
sopiga_recipient_detail_id BIGINT,
sopiga_external_id VARCHAR(100),
created_at TIMESTAMP DEFAULT NOW(),
updated_at TIMESTAMP DEFAULT NOW(),
dispatched_at TIMESTAMP,
delivered_at TIMESTAMP,
failed_at TIMESTAMP,
CONSTRAINT status_valid CHECK (status IN ('pending', 'dispatched', 'delivered', 'failed', 'retry_scheduled')),
CONSTRAINT error_count_positive CHECK (error_count >= 0)
);
CREATE INDEX idx_broadcast_staging_status ON collection_broadcasts.broadcast_staging(status);
CREATE INDEX idx_broadcast_staging_created_at ON collection_broadcasts.broadcast_staging(created_at DESC);
CREATE INDEX idx_broadcast_staging_collar_id ON collection_broadcasts.broadcast_staging(sopiga_collar_id);
CREATE INDEX idx_broadcast_staging_template_id ON collection_broadcasts.broadcast_staging(sopiga_template_id);
CREATE INDEX idx_broadcast_staging_pending_query
ON collection_broadcasts.broadcast_staging(status, created_at ASC)
WHERE status = 'pending';
CREATE INDEX idx_broadcast_staging_retry_query
ON collection_broadcasts.broadcast_staging(status, error_count, created_at ASC)
WHERE status = 'retry_scheduled' AND error_count < 3;
-- ============================================================================
-- 3. AUDIT & LOGGING
-- ============================================================================
CREATE TABLE collection_broadcasts.broadcast_audit_log (
id BIGSERIAL PRIMARY KEY,
broadcast_id BIGINT NOT NULL REFERENCES collection_broadcasts.broadcast_staging(id) ON DELETE CASCADE,
old_status VARCHAR(50),
new_status VARCHAR(50) NOT NULL,
reason VARCHAR(500),
sopiga_response JSONB,
changed_by VARCHAR(100) DEFAULT 'system',
created_at TIMESTAMP DEFAULT NOW()
);
CREATE INDEX idx_broadcast_audit_log_broadcast_id ON collection_broadcasts.broadcast_audit_log(broadcast_id);
CREATE INDEX idx_broadcast_audit_log_created_at ON collection_broadcasts.broadcast_audit_log(created_at DESC);
CREATE TABLE collection_broadcasts.broadcast_error_log (
id BIGSERIAL PRIMARY KEY,
broadcast_id BIGINT NOT NULL REFERENCES collection_broadcasts.broadcast_staging(id) ON DELETE CASCADE,
error_type VARCHAR(100),
error_code VARCHAR(50),
error_message TEXT,
error_details JSONB,
attempt_number INT,
next_retry_at TIMESTAMP,
created_at TIMESTAMP DEFAULT NOW()
);
CREATE INDEX idx_broadcast_error_log_broadcast_id ON collection_broadcasts.broadcast_error_log(broadcast_id);
CREATE INDEX idx_broadcast_error_log_error_type ON collection_broadcasts.broadcast_error_log(error_type);
CREATE TABLE collection_broadcasts.sopiga_sync_job (
id BIGSERIAL PRIMARY KEY,
broadcast_id BIGINT NOT NULL REFERENCES collection_broadcasts.broadcast_staging(id) ON DELETE CASCADE,
sopiga_recipient_detail_id BIGINT,
last_synced_at TIMESTAMP,
last_status_from_sopiga VARCHAR(50),
sync_count INT DEFAULT 0,
next_sync_at TIMESTAMP,
completed_at TIMESTAMP,
created_at TIMESTAMP DEFAULT NOW()
);
CREATE INDEX idx_sopiga_sync_job_broadcast_id ON collection_broadcasts.sopiga_sync_job(broadcast_id);
CREATE INDEX idx_sopiga_sync_job_next_sync_at ON collection_broadcasts.sopiga_sync_job(next_sync_at);
-- ============================================================================
-- 4. VIEWS
-- ============================================================================
CREATE VIEW collection_broadcasts.v_template_variables_ordered AS
SELECT
tvm.sopiga_template_id,
stc.template_name,
tvm.variable_order,
tvm.sopiga_variable_name,
tvm.variable_type,
tvm.db_field_source,
tvm.is_required,
tvm.example_value
FROM collection_broadcasts.template_variable_mapping tvm
JOIN collection_broadcasts.sopiga_template_config stc ON tvm.sopiga_template_id = stc.sopiga_template_id
WHERE stc.active = TRUE
ORDER BY tvm.sopiga_template_id, tvm.variable_order;
CREATE VIEW collection_broadcasts.v_collar_summary AS
SELECT
sc.sopiga_collar_id,
sc.broadcast_name,
COUNT(*) as total_records,
COUNT(CASE WHEN bs.status = 'pending' THEN 1 END) as pending,
COUNT(CASE WHEN bs.status = 'dispatched' THEN 1 END) as dispatched,
COUNT(CASE WHEN bs.status = 'delivered' THEN 1 END) as delivered,
COUNT(CASE WHEN bs.status = 'failed' THEN 1 END) as failed,
ROUND(100.0 * COUNT(CASE WHEN bs.status = 'delivered' THEN 1 END) / NULLIF(COUNT(*), 0), 2) as delivery_rate_percent
FROM collection_broadcasts.sopiga_collar_config sc
LEFT JOIN collection_broadcasts.broadcast_staging bs ON sc.sopiga_collar_id = bs.sopiga_collar_id
GROUP BY sc.sopiga_collar_id, sc.broadcast_name;
CREATE VIEW collection_broadcasts.v_failed_records_24h AS
SELECT
id,
sopiga_collar_id,
sopiga_template_id,
message_payload->>'nasabah_nama' as nasabah_nama,
message_payload->>'nasabah_phone' as nasabah_phone,
error_message,
error_count,
failed_at,
created_at
FROM collection_broadcasts.broadcast_staging
WHERE status = 'failed' AND created_at > NOW() - INTERVAL '1 DAY'
ORDER BY failed_at DESC;
CREATE VIEW collection_broadcasts.v_delivery_rate_24h AS
SELECT
COUNT(*) as total,
COUNT(CASE WHEN status = 'delivered' THEN 1 END) as delivered,
COUNT(CASE WHEN status = 'failed' THEN 1 END) as failed,
COUNT(CASE WHEN status = 'dispatched' THEN 1 END) as in_progress,
ROUND(100.0 * COUNT(CASE WHEN status = 'delivered' THEN 1 END) / NULLIF(COUNT(*), 0), 2) as delivery_rate_percent
FROM collection_broadcasts.broadcast_staging
WHERE created_at > NOW() - INTERVAL '1 DAY';
-- ============================================================================
-- 5. HELPER FUNCTIONS
-- ============================================================================
CREATE OR REPLACE FUNCTION collection_broadcasts.get_template_variables(
p_sopiga_template_id INT
)
RETURNS TABLE(
variable_order INT,
sopiga_variable_name VARCHAR,
variable_type VARCHAR,
db_field_source VARCHAR,
is_required BOOLEAN
) AS $$
BEGIN
RETURN QUERY
SELECT
tvm.variable_order,
tvm.sopiga_variable_name,
tvm.variable_type,
tvm.db_field_source,
tvm.is_required
FROM collection_broadcasts.template_variable_mapping tvm
WHERE tvm.sopiga_template_id = p_sopiga_template_id
ORDER BY tvm.variable_order ASC;
END;
$$ LANGUAGE plpgsql;
CREATE OR REPLACE FUNCTION collection_broadcasts.update_broadcast_status(
p_broadcast_id BIGINT,
p_new_status VARCHAR,
p_error_message TEXT DEFAULT NULL,
p_sopiga_response JSONB DEFAULT NULL
)
RETURNS VOID AS $$
DECLARE
v_old_status VARCHAR;
BEGIN
SELECT status INTO v_old_status
FROM collection_broadcasts.broadcast_staging
WHERE id = p_broadcast_id;
UPDATE collection_broadcasts.broadcast_staging
SET
status = p_new_status,
error_message = p_error_message,
updated_at = NOW(),
dispatched_at = CASE WHEN p_new_status = 'dispatched' THEN NOW() ELSE dispatched_at END,
delivered_at = CASE WHEN p_new_status = 'delivered' THEN NOW() ELSE delivered_at END,
failed_at = CASE WHEN p_new_status = 'failed' THEN NOW() ELSE failed_at END
WHERE id = p_broadcast_id;
INSERT INTO collection_broadcasts.broadcast_audit_log (broadcast_id, old_status, new_status, sopiga_response)
VALUES (p_broadcast_id, v_old_status, p_new_status, p_sopiga_response);
END;
$$ LANGUAGE plpgsql;