111 lines
3.5 KiB
Go
111 lines
3.5 KiB
Go
package store
|
|
|
|
import (
|
|
"database/sql"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"time"
|
|
|
|
"git.ipao.vip/rogee/go-sip/internal/contract"
|
|
)
|
|
|
|
var ErrQuerySnapshotTooLarge = errors.New("query snapshot exceeds MQ message budget")
|
|
|
|
func callQuerySnapshot(tx *sql.Tx, request contract.ServiceMessage, callID string, now time.Time) (map[string]any, error) {
|
|
rows, err := tx.Query(`SELECT body,status FROM outbox WHERE tenant_key=?
|
|
AND json_extract(body,'$.tenant_id')=? AND json_extract(body,'$.dispatcher_id')=?
|
|
AND json_extract(body,'$.event_type') IS NOT NULL
|
|
AND json_extract(body,'$.payload.call_id')=?
|
|
ORDER BY json_extract(body,'$.occurred_at'),event_id`, request.TenantKey, request.TenantID, request.DispatcherID, callID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
attempts, transcripts, recordings := []json.RawMessage{}, []json.RawMessage{}, []json.RawMessage{}
|
|
delivery := map[string]int{"pending": 0, "retry": 0, "dispatching": 0, "published": 0}
|
|
var current map[string]any
|
|
var currentVersion int64 = -1
|
|
var finished map[string]any
|
|
var finishedVersion int64 = -1
|
|
count, size := 0, 0
|
|
for rows.Next() {
|
|
var raw []byte
|
|
var status string
|
|
if err := rows.Scan(&raw, &status); err != nil {
|
|
return nil, err
|
|
}
|
|
if _, ok := delivery[status]; !ok {
|
|
return nil, fmt.Errorf("unknown outbox delivery state %q", status)
|
|
}
|
|
delivery[status]++
|
|
count++
|
|
var event struct {
|
|
EventType string `json:"event_type"`
|
|
Payload json.RawMessage `json:"payload"`
|
|
}
|
|
if err := json.Unmarshal(raw, &event); err != nil {
|
|
return nil, err
|
|
}
|
|
switch event.EventType {
|
|
case "call.status", "call.finished":
|
|
var fields map[string]any
|
|
if err := json.Unmarshal(event.Payload, &fields); err != nil {
|
|
return nil, err
|
|
}
|
|
var state struct {
|
|
Version int64 `json:"call_version"`
|
|
}
|
|
if err := json.Unmarshal(event.Payload, &state); err != nil {
|
|
return nil, err
|
|
}
|
|
if state.Version < 1 {
|
|
return nil, errors.New("call fact has no positive call_version")
|
|
}
|
|
if event.EventType == "call.status" {
|
|
attempts = append(attempts, json.RawMessage(raw))
|
|
size += len(raw)
|
|
if state.Version >= currentVersion {
|
|
current, currentVersion = fields, state.Version
|
|
}
|
|
} else if state.Version >= finishedVersion {
|
|
finished, finishedVersion = fields, state.Version
|
|
}
|
|
case "transcript.updated", "transcript.failed":
|
|
transcripts = append(transcripts, json.RawMessage(raw))
|
|
size += len(raw)
|
|
case "recording.uploaded", "recording.failed":
|
|
recordings = append(recordings, json.RawMessage(raw))
|
|
size += len(raw)
|
|
}
|
|
if size > 256<<10 {
|
|
return nil, ErrQuerySnapshotTooLarge
|
|
}
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
if count == 0 {
|
|
return nil, sql.ErrNoRows
|
|
}
|
|
if finished != nil && finishedVersion >= currentVersion {
|
|
current, currentVersion = finished, finishedVersion
|
|
// A persisted call.finished fact establishes that the call has ended.
|
|
current["call_state"] = "ended"
|
|
}
|
|
if current == nil {
|
|
return nil, errors.New("call facts contain no state/version evidence")
|
|
}
|
|
result := map[string]any{"call_id": callID, "call_version": currentVersion, "snapshot_at": now.Format(time.RFC3339Nano)}
|
|
result["attempts"] = attempts
|
|
result["transcript"] = map[string]any{"events": transcripts}
|
|
result["recordings"] = recordings
|
|
result["delivery"] = delivery
|
|
for _, key := range []string{"execution_id", "task_id", "task_item_id", "call_state", "reason_code", "outcome", "started_at", "ended_at", "duration_ms"} {
|
|
if value, ok := current[key]; ok {
|
|
result[key] = value
|
|
}
|
|
}
|
|
return result, nil
|
|
}
|