Files
go-sip/internal/store/current_upload.go
T

138 lines
7.2 KiB
Go

package store
import (
"database/sql"
"encoding/hex"
"errors"
"fmt"
"strings"
)
var ErrCurrentUploadConflict = errors.New("recording upload differs from its original approved target")
var ErrCurrentUploadNotFound = errors.New("no durable original recording target")
var ErrCurrentUploadAlreadyConfirmed = errors.New("recording upload already confirmed; another PUT is forbidden")
// CurrentUploadProof is the authenticated Agent's successful PUT observation;
// the Dispatcher also verifies every asset field against its durable grant.
// It does not claim an independent OSS HEAD or SaaS application receipt.
type CurrentUploadProof struct {
UploadID string
RecordingID string
StatusCode int
SizeBytes int64
SHA256 string
}
// CurrentRecordingGrant holds the original D-approved OSS destination. Signed
// URLs, headers and short-lived tokens are deliberately never stored here.
type CurrentRecordingGrant struct {
DispatcherID string
SourceEventID string
UploadID string
RecordingID string
Bucket string
ObjectKey string
ChecksumSHA256 string
SizeBytes int64
Format string
Channels int
SampleRateHz int
DurationMS int64
ConfirmedAt string
}
// RequireReservedCall ties Agent-reported facts to the numeric tenant and the
// durable, already reserved execution. Callers must also authenticate the
// active Agent session before accepting any report or signing an OSS token.
func (s *CurrentStore) RequireReservedCall(dispatcherID, sourceEventID string, tenantID int64) error {
if dispatcherID == "" || sourceEventID == "" || tenantID <= 0 {
return errors.New("recording source lacks a Dispatcher execution identity")
}
var storedTenant int64
var status, trunkID string
var snapshot []byte
if err := s.db.QueryRow(`SELECT tenant_id,status,COALESCE(selected_trunk_id,''),snapshot_json FROM dispatcher_inbox WHERE dispatcher_id=? AND event_id=?`, dispatcherID, sourceEventID).Scan(&storedTenant, &status, &trunkID, &snapshot); err != nil {
return fmt.Errorf("load reserved call identity: %w", err)
}
if storedTenant != tenantID || trunkID == "" || len(snapshot) == 0 || (status != "dispatching" && status != "dispatched" && status != "unknown" && status != "finished") {
return errors.New("recording source is not the approved tenant execution")
}
return nil
}
func (s *CurrentStore) BindRecordingUpload(proposed CurrentRecordingGrant) (CurrentRecordingGrant, bool, error) {
if proposed.DispatcherID == "" || proposed.SourceEventID == "" || proposed.UploadID == "" || proposed.RecordingID == "" || proposed.Bucket == "" || proposed.ObjectKey == "" || proposed.SizeBytes <= 0 || proposed.Format != "wav" || (proposed.Channels != 1 && proposed.Channels != 2) || proposed.SampleRateHz != 16000 || proposed.DurationMS < 0 || proposed.ConfirmedAt != "" {
return CurrentRecordingGrant{}, false, errors.New("recording grant lacks a bounded original asset and execution")
}
if len(proposed.ChecksumSHA256) != 64 || strings.ToLower(proposed.ChecksumSHA256) != proposed.ChecksumSHA256 {
return CurrentRecordingGrant{}, false, errors.New("recording grant requires a canonical SHA-256 digest")
}
if _, err := hex.DecodeString(proposed.ChecksumSHA256); err != nil {
return CurrentRecordingGrant{}, false, errors.New("recording grant has an invalid SHA-256 digest")
}
tx, err := s.db.Begin()
if err != nil {
return CurrentRecordingGrant{}, false, err
}
defer tx.Rollback()
var status, trunkID string
var snapshot []byte
err = tx.QueryRow(`SELECT status,COALESCE(selected_trunk_id,''),snapshot_json FROM dispatcher_inbox WHERE dispatcher_id=? AND event_id=?`, proposed.DispatcherID, proposed.SourceEventID).Scan(&status, &trunkID, &snapshot)
if errors.Is(err, sql.ErrNoRows) {
return CurrentRecordingGrant{}, false, ErrCurrentUploadNotFound
}
if err != nil {
return CurrentRecordingGrant{}, false, fmt.Errorf("load approved call for recording: %w", err)
}
if status != "dispatching" && status != "dispatched" && status != "unknown" && status != "finished" || trunkID == "" || len(snapshot) == 0 {
return CurrentRecordingGrant{}, false, errors.New("recording grant requires an already reserved approved call")
}
inserted, err := tx.Exec(`INSERT INTO dispatcher_recordings(dispatcher_id,source_event_id,upload_id,recording_id,bucket,object_key,checksum_sha256,size_bytes,format,channels,sample_rate_hz,duration_ms)
VALUES(?,?,?,?,?,?,?,?,?,?,?,?) ON CONFLICT(dispatcher_id,source_event_id) DO NOTHING`, proposed.DispatcherID, proposed.SourceEventID, proposed.UploadID, proposed.RecordingID, proposed.Bucket, proposed.ObjectKey, proposed.ChecksumSHA256, proposed.SizeBytes, proposed.Format, proposed.Channels, proposed.SampleRateHz, proposed.DurationMS)
if err != nil {
return CurrentRecordingGrant{}, false, fmt.Errorf("persist immutable recording target: %w", err)
}
count, err := inserted.RowsAffected()
if err != nil {
return CurrentRecordingGrant{}, false, err
}
stored, err := loadRecordingUpload(tx.QueryRow(`SELECT dispatcher_id,source_event_id,upload_id,recording_id,bucket,object_key,checksum_sha256,size_bytes,format,channels,sample_rate_hz,duration_ms,COALESCE(confirmed_at,'') FROM dispatcher_recordings WHERE dispatcher_id=? AND source_event_id=?`, proposed.DispatcherID, proposed.SourceEventID))
if err != nil {
return CurrentRecordingGrant{}, false, fmt.Errorf("inspect persisted recording target: %w", err)
}
if stored.ConfirmedAt != "" {
return CurrentRecordingGrant{}, false, ErrCurrentUploadAlreadyConfirmed
}
if stored != proposed {
return CurrentRecordingGrant{}, false, ErrCurrentUploadConflict
}
var finalized int
err = tx.QueryRow(`SELECT 1 FROM dispatcher_outbox WHERE dispatcher_id=? AND event_id=?`, proposed.DispatcherID, currentResultEventID(proposed.DispatcherID, proposed.SourceEventID)).Scan(&finalized)
if err == nil {
return CurrentRecordingGrant{}, false, ErrCurrentResultConflict
}
if !errors.Is(err, sql.ErrNoRows) {
return CurrentRecordingGrant{}, false, fmt.Errorf("check existing call result before granting an upload: %w", err)
}
if err := tx.Commit(); err != nil {
return CurrentRecordingGrant{}, false, fmt.Errorf("commit original recording target: %w", err)
}
return stored, count == 1, nil
}
func (s *CurrentStore) LoadRecordingUpload(dispatcherID, sourceEventID string) (CurrentRecordingGrant, error) {
if dispatcherID == "" || sourceEventID == "" {
return CurrentRecordingGrant{}, ErrCurrentUploadNotFound
}
return loadRecordingUpload(s.db.QueryRow(`SELECT dispatcher_id,source_event_id,upload_id,recording_id,bucket,object_key,checksum_sha256,size_bytes,format,channels,sample_rate_hz,duration_ms,COALESCE(confirmed_at,'') FROM dispatcher_recordings WHERE dispatcher_id=? AND source_event_id=?`, dispatcherID, sourceEventID))
}
func loadRecordingUpload(row *sql.Row) (CurrentRecordingGrant, error) {
var grant CurrentRecordingGrant
err := row.Scan(&grant.DispatcherID, &grant.SourceEventID, &grant.UploadID, &grant.RecordingID, &grant.Bucket, &grant.ObjectKey, &grant.ChecksumSHA256, &grant.SizeBytes, &grant.Format, &grant.Channels, &grant.SampleRateHz, &grant.DurationMS, &grant.ConfirmedAt)
if errors.Is(err, sql.ErrNoRows) {
return CurrentRecordingGrant{}, ErrCurrentUploadNotFound
}
return grant, err
}