90 lines
3.3 KiB
Go
90 lines
3.3 KiB
Go
package bell_connector
|
|
|
|
import (
|
|
"context"
|
|
"crypto/tls"
|
|
"errors"
|
|
"fmt"
|
|
"log"
|
|
"os"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/gin-gonic/gin"
|
|
"gorm.io/gorm"
|
|
|
|
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/integration/machine_identity"
|
|
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/outbox"
|
|
)
|
|
|
|
type Environment func(string) string
|
|
|
|
func StartRuntime(ctx context.Context, engine *gin.Engine, db *gorm.DB, getenv Environment) error {
|
|
if getenv == nil {
|
|
getenv = os.Getenv
|
|
}
|
|
ingressEnabled := enabled(getenv("SENSE_EVENT_INGRESS_ENABLED"))
|
|
relayEnabled := enabled(getenv("SENSE_BELL_CONNECTOR_ENABLED"))
|
|
if !ingressEnabled && !relayEnabled {
|
|
return nil
|
|
}
|
|
if engine == nil || db == nil {
|
|
return errors.New("Sense connector runtime requires engine and database")
|
|
}
|
|
for _, model := range []any{&InboundEvent{}, &EvidenceRecord{}, &ReplayToken{}, &outbox.Message{}, &outbox.DeliveryRecord{}, &outbox.Attempt{}} {
|
|
if !db.Migrator().HasTable(model) {
|
|
return fmt.Errorf("Sense connector migration is not applied for %T", model)
|
|
}
|
|
}
|
|
if ingressEnabled {
|
|
registry, err := machine_identity.LoadRegistry(getenv("SENSE_MACHINE_PRINCIPAL_REGISTRY"), "yovision-sense")
|
|
if err != nil {
|
|
return fmt.Errorf("load Sense machine identity registry: %w", err)
|
|
}
|
|
verifier := machine_identity.Verifier{Registry: registry, Replay: PersistentReplayStore{DB: db}}
|
|
engine.POST("/v1/events", (IngressHandler{DB: db, Verifier: verifier, EvidenceOwnerID: strings.TrimSpace(getenv("SENSE_EVIDENCE_OWNER_ID"))}).Post)
|
|
engine.GET("/v1/evidence/:evidence_id", (EvidenceHandler{DB: db, Verifier: verifier}).Get)
|
|
}
|
|
if !relayEnabled {
|
|
return nil
|
|
}
|
|
privateKey, err := machine_identity.LoadPrivateKey(getenv("SENSE_BELL_PRIVATE_KEY_PATH"))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
signer := machine_identity.Signer{Principal: strings.TrimSpace(getenv("SENSE_BELL_PRINCIPAL_ID")), KeyID: strings.TrimSpace(getenv("SENSE_BELL_KEY_ID")), PrivateKey: privateKey}
|
|
policy := machine_identity.TransportPolicy{TLSMinVersion: tls.VersionTLS12, VerifyCertificate: true, VerifyHostname: true, ConnectTimeout: 5 * time.Second, ResponseHeaderTimeout: 10 * time.Second, RequestTimeout: 15 * time.Second, MaxRequestBytes: MaxInboundBytes}
|
|
client, err := NewClient(strings.TrimSpace(getenv("SENSE_BELL_ENDPOINT")), strings.TrimSpace(getenv("SENSE_RELAY_ID")), signer, policy)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
interval := 2 * time.Second
|
|
if raw := strings.TrimSpace(getenv("SENSE_BELL_RELAY_INTERVAL_MS")); raw != "" {
|
|
milliseconds, parseErr := strconv.Atoi(raw)
|
|
if parseErr != nil || milliseconds < 100 || milliseconds > 60000 {
|
|
return errors.New("SENSE_BELL_RELAY_INTERVAL_MS must be between 100 and 60000")
|
|
}
|
|
interval = time.Duration(milliseconds) * time.Millisecond
|
|
}
|
|
go runRelay(ctx, Relay{DB: db, Client: client}, interval)
|
|
return nil
|
|
}
|
|
|
|
func runRelay(ctx context.Context, relay Relay, interval time.Duration) {
|
|
ticker := time.NewTicker(interval)
|
|
defer ticker.Stop()
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-ticker.C:
|
|
if _, err := relay.DeliverBatch(ctx, "sense-bell-runtime", 50); err != nil && ctx.Err() == nil {
|
|
log.Printf("Sense Bell connector delivery failed: %v", err)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
func enabled(value string) bool { return strings.EqualFold(strings.TrimSpace(value), "true") }
|