| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617 |
- package main
- import (
- "context"
- "encoding/json"
- "errors"
- "fmt"
- "log/slog"
- "net/http"
- "os"
- "os/signal"
- "strings"
- "syscall"
- "time"
- "vocat/internal/auth"
- "vocat/internal/config"
- "vocat/internal/device"
- "vocat/internal/loghub"
- "vocat/internal/server"
- "vocat/internal/store"
- "vocat/internal/update"
- "vocat/internal/vowifi"
- "vocat/internal/vowifi/ike"
- "vocat/internal/vowifi/ims"
- "vocat/internal/vowifi/integration"
- vowifiruntime "vocat/internal/vowifi/runtime"
- "vocat/web"
- )
- func main() {
- logs := loghub.New(slog.NewJSONHandler(os.Stdout, nil), 2000)
- logger := slog.New(logs)
- args := os.Args[1:]
- switch subcommand, rest := splitSubcommand(args); subcommand {
- case "":
- // No subcommand: run the server. Backward-compatible with the
- // existing systemd unit (ExecStart=/opt/vocat/bin/vocat).
- if err := run(logger, logs); err != nil {
- logger.Error("server stopped", "error", err)
- os.Exit(1)
- }
- case "version", "-v", "--version":
- runVersion()
- case "update":
- if err := update.Run(logger, rest); err != nil {
- logger.Error("update failed", "error", err)
- os.Exit(1)
- }
- case "menu":
- if err := runMenu(logger); err != nil {
- logger.Error("menu failed", "error", err)
- os.Exit(1)
- }
- case "help", "-h", "--help":
- printUsage(os.Stdout)
- default:
- fmt.Fprintf(os.Stderr, "vocat: unknown subcommand %q\n\n", subcommand)
- printUsage(os.Stderr)
- os.Exit(2)
- }
- }
- // splitSubcommand returns the first non-flag token as the subcommand and the
- // remaining args. An empty arg list yields ("", nil) → server mode.
- func splitSubcommand(args []string) (string, []string) {
- if len(args) == 0 {
- return "", nil
- }
- return args[0], args[1:]
- }
- func run(logger *slog.Logger, logs *loghub.Hub) error {
- cfg, err := config.Load()
- if err != nil {
- return fmt.Errorf("load configuration: %w", err)
- }
- if cfg.UsesDefaultCredentials() {
- logger.Warn(
- "default admin credentials are active; set VOCAT_ADMIN_PASSWORD before exposing the service",
- )
- }
- startupContext, cancelStartup := context.WithTimeout(context.Background(), 15*time.Second)
- defer cancelStartup()
- database, err := store.Open(startupContext, cfg.DatabasePath)
- if err != nil {
- return err
- }
- defer database.Close()
- authService, err := auth.New(database, auth.Options{
- SessionTTL: cfg.SessionTTL,
- })
- if err != nil {
- return err
- }
- if err := authService.EnsureAdmin(
- startupContext,
- cfg.AdminUsername,
- cfg.AdminPassword,
- ); err != nil {
- return err
- }
- deviceManager, err := device.NewManager(device.Options{})
- if err != nil {
- return fmt.Errorf("create device manager: %w", err)
- }
- if err := deviceManager.Start(startupContext); err != nil {
- logger.Warn("device discovery is not available at startup", "error", err)
- }
- if err := provisionDiscoveredDevices(startupContext, database, deviceManager); err != nil {
- logger.Warn("automatic first-run device provisioning failed", "error", err)
- }
- defer func() {
- stopContext, cancel := context.WithTimeout(context.Background(), 5*time.Second)
- defer cancel()
- if err := deviceManager.Stop(stopContext); err != nil {
- logger.Warn("stop device manager", "error", err)
- }
- }()
- pollContext, cancelPolling := context.WithCancel(context.Background())
- defer cancelPolling()
- go pollDeviceSnapshots(pollContext, logger, database, deviceManager)
- go persistLogsToStore(pollContext, logger, logs, database)
- vowifiManager, err := configureVoWiFiRuntime(
- startupContext,
- logger,
- database,
- deviceManager,
- )
- if err != nil {
- return fmt.Errorf("configure VoWiFi runtime: %w", err)
- }
- defer func() {
- stopContext, cancel := context.WithTimeout(context.Background(), 15*time.Second)
- defer cancel()
- if err := vowifiManager.Close(stopContext); err != nil {
- logger.Warn("stop VoWiFi runtime", "error", err)
- }
- }()
- handler, err := server.New(server.Options{
- Store: database,
- Auth: authService,
- Devices: deviceManager,
- VoWiFi: vowifiManager,
- Logs: logs,
- Assets: web.Dist,
- Logger: logger,
- SecureCookies: cfg.SecureCookies,
- MaxRequestBodyBytes: cfg.MaxRequestBodyBytes,
- })
- if err != nil {
- return err
- }
- go handler.StartLogRetentionLoop(pollContext, time.Minute)
- go handler.StartSMSSyncLoop(pollContext, 15*time.Second)
- handler.StartTelegramBot(pollContext)
- handler.StartSMSNotificationDispatchers(pollContext)
- httpServer := &http.Server{
- Addr: cfg.Address,
- Handler: handler,
- ReadHeaderTimeout: 5 * time.Second,
- ReadTimeout: 15 * time.Second,
- WriteTimeout: 30 * time.Second,
- IdleTimeout: 90 * time.Second,
- MaxHeaderBytes: 1 << 20,
- }
- signalContext, stopSignals := signal.NotifyContext(
- context.Background(),
- os.Interrupt,
- syscall.SIGTERM,
- )
- defer stopSignals()
- serverError := make(chan error, 1)
- go func() {
- logger.Info("HTTP server listening", "address", cfg.Address)
- err := httpServer.ListenAndServe()
- if errors.Is(err, http.ErrServerClosed) {
- err = nil
- }
- serverError <- err
- }()
- select {
- case err := <-serverError:
- return err
- case <-signalContext.Done():
- logger.Info("shutdown signal received")
- }
- shutdownContext, cancelShutdown := context.WithTimeout(
- context.Background(),
- cfg.ShutdownTimeout,
- )
- defer cancelShutdown()
- if err := httpServer.Shutdown(shutdownContext); err != nil {
- _ = httpServer.Close()
- return fmt.Errorf("graceful HTTP shutdown: %w", err)
- }
- return <-serverError
- }
- func configureVoWiFiRuntime(
- ctx context.Context,
- logger *slog.Logger,
- database *store.Store,
- deviceManager *device.Manager,
- ) (*vowifiruntime.Manager, error) {
- mapper := integration.ATMapper{
- Store: database,
- Devices: deviceManager,
- }
- adapter, err := vowifi.NewEC20Adapter(mapper, vowifi.EC20AdapterOptions{
- // The test deployment is deliberately non-cellular. VoWiFi teardown
- // may restore CFUN, but it must never reactivate a PDP context.
- RestoreCellularData: false,
- })
- if err != nil {
- return nil, err
- }
- projector := integration.StateProjector{
- Store: database,
- Devices: mapper,
- }
- manager := vowifiruntime.New(vowifiruntime.Options{
- Logger: logger,
- OnState: projector.Save,
- Factory: func(factoryContext context.Context, deviceID string) (*vowifi.Orchestrator, error) {
- deviceConfig, err := database.Device(factoryContext, deviceID)
- if err != nil {
- return nil, fmt.Errorf("load device %q VoWiFi config: %w", deviceID, err)
- }
- return newVoWiFiOrchestrator(deviceConfig, database, adapter)
- },
- })
- configured, err := database.ListDevices(ctx)
- if err != nil {
- _ = manager.Close(context.Background())
- return nil, err
- }
- for _, deviceConfig := range configured {
- if err := manager.Ensure(ctx, deviceConfig.ID); err != nil {
- _ = manager.Close(context.Background())
- return nil, fmt.Errorf("register device %q VoWiFi runtime: %w", deviceConfig.ID, err)
- }
- if deviceConfig.VoWiFiEnabled {
- if _, err := manager.RequestEnabled(deviceConfig.ID, true); err != nil {
- _ = manager.Close(context.Background())
- return nil, fmt.Errorf("start device %q VoWiFi policy: %w", deviceConfig.ID, err)
- }
- }
- }
- return manager, nil
- }
- func newVoWiFiOrchestrator(
- deviceConfig store.Device,
- database *store.Store,
- adapter *vowifi.EC20Adapter,
- ) (*vowifi.Orchestrator, error) {
- apn := deviceConfig.APN
- if apn == "" {
- apn = "ims"
- }
- tunnelProvider, err := ike.NewProvider(ike.Config{APN: apn})
- if err != nil {
- return nil, fmt.Errorf("device %q IKE provider: %w", deviceConfig.ID, err)
- }
- imsProvider, err := ims.NewProvider(adapter, ims.Config{
- // The userspace SWu data plane currently carries the protected P-CSCF
- // signalling path over TCP.
- Transport: "tcp",
- // Some Vodafone UK SIM profiles leave AT+CSCA empty; Vodafone publishes
- // this service-centre number for manual SMS setup.
- SMSCenter: "+447785016005",
- OnSMS: func(ctx context.Context, message ims.ReceivedSMS) error {
- extra, _ := json.Marshal(map[string]any{
- "transport": "ims",
- "encoding": message.Encoding,
- "concat": message.Concat,
- "rp_reference": message.RPReference,
- "call_id": message.CallID,
- "received_at": message.Timestamp,
- "service_center_timestamp": message.ServiceCenterTimestamp,
- "raw_rpdu": message.RawRPDU,
- "raw_tpdu": message.RawTPDU,
- })
- partsTotal := 1
- if message.Concat != nil && message.Concat.Total > 0 {
- partsTotal = message.Concat.Total
- }
- _, saveErr := database.SaveSMSMessage(ctx, store.SMSMessage{
- MessageID: message.MessageID,
- DeviceID: message.DeviceID,
- IMSI: message.IMSI,
- Peer: message.From,
- Direction: "inbound",
- Body: message.Text,
- Timestamp: message.Timestamp,
- Status: "received",
- Source: "ims",
- PartsTotal: partsTotal,
- Read: false,
- Extra: extra,
- })
- return saveErr
- },
- OnSMSStatus: func(ctx context.Context, report ims.ReceivedSMSStatus) error {
- deliveryReport := store.SMSDeliveryReport{
- DeviceID: report.DeviceID,
- IMSI: report.IMSI,
- Peer: report.To,
- Source: "ims",
- MessageReference: report.MessageReference,
- StatusCode: report.StatusCode,
- DeliveryState: report.DeliveryStatus,
- ServiceCenterTime: report.ServiceCenterTimestamp,
- DischargeTime: report.DischargeTimestamp,
- ReceivedAt: report.Timestamp,
- }
- var applyErr error
- for attempt := 0; attempt < 10; attempt++ {
- _, applyErr = database.ApplySMSDeliveryReport(ctx, deliveryReport)
- if !errors.Is(applyErr, store.ErrNotFound) {
- return applyErr
- }
- // A status report can race the API handler persisting the SIP 202
- // result. Give that write a brief chance to complete.
- select {
- case <-ctx.Done():
- return ctx.Err()
- case <-time.After(100 * time.Millisecond):
- }
- }
- // A late report from before this process started must still be
- // acknowledged, otherwise the SMSC will keep retransmitting it.
- return nil
- },
- })
- if err != nil {
- return nil, fmt.Errorf("device %q IMS provider: %w", deviceConfig.ID, err)
- }
- orchestrator, err := vowifi.New(vowifi.Dependencies{
- SIM: adapter,
- AKA: adapter,
- Radio: adapter,
- Proxy: integration.ProxyResolver{Store: database},
- Tunnel: tunnelProvider,
- IMS: imsProvider,
- Phones: integration.PhoneStore{Store: database, DeviceID: deviceConfig.ID},
- }, vowifi.Options{
- DeviceID: deviceConfig.ID,
- AllowIMSWithoutSMS: true,
- })
- if err != nil {
- return nil, fmt.Errorf("device %q VoWiFi orchestrator: %w", deviceConfig.ID, err)
- }
- return orchestrator, nil
- }
- func provisionDiscoveredDevices(
- ctx context.Context,
- database *store.Store,
- manager *device.Manager,
- ) error {
- configured, err := database.ListDevices(ctx)
- if err != nil {
- return err
- }
- if len(configured) != 0 {
- return nil
- }
- for _, discovered := range manager.List() {
- candidate := discovered.Candidate
- backend := "at"
- control := candidate.ATPort.OpenPath()
- if candidate.QMIControl != "" {
- backend = "qmi"
- control = candidate.QMIControl
- }
- name := candidate.Product
- if name == "" || strings.EqualFold(name, "Android") {
- name = "Quectel EC20 / EC25"
- }
- if err := database.UpsertDevice(ctx, store.Device{
- ID: discovered.ID,
- Name: name,
- Interface: candidate.NetworkInterface,
- ControlDevice: control,
- ATPort: candidate.ATPort.OpenPath(),
- USBPath: candidate.USBPath,
- ProxyPort: 1080,
- BaudRate: 115200,
- DataBits: 8,
- StopBits: 1,
- Parity: "none",
- DeviceBackend: backend,
- ESIMTransport: backend,
- NetworkEnabled: false,
- SMSEnabled: true,
- VoWiFiEnabled: false,
- }); err != nil {
- return err
- }
- }
- return nil
- }
- // persistLogsToStore subscribes to the live log hub and durably appends every
- // entry to the log_events table, so runtime logs survive restarts and can be
- // pruned by the configured retention policy.
- func persistLogsToStore(
- ctx context.Context,
- logger *slog.Logger,
- logs *loghub.Hub,
- database *store.Store,
- ) {
- entries, cancel := logs.Subscribe(512)
- defer cancel()
- for {
- select {
- case <-ctx.Done():
- return
- case entry, ok := <-entries:
- if !ok {
- return
- }
- var fields json.RawMessage
- if len(entry.Fields) > 0 {
- if raw, err := json.Marshal(entry.Fields); err == nil {
- fields = raw
- }
- }
- if _, err := database.AppendLogEvent(ctx, store.LogEvent{
- Time: entry.Time,
- Level: entry.Level,
- Message: entry.Message,
- Caller: entry.Caller,
- Fields: fields,
- }); err != nil && ctx.Err() == nil {
- logger.Warn("persist log event failed", "error", err)
- }
- }
- }
- }
- func pollDeviceSnapshots(
- ctx context.Context,
- logger *slog.Logger,
- database *store.Store,
- manager *device.Manager,
- ) {
- refresh := func() {
- discoveryContext, cancelDiscovery := context.WithTimeout(ctx, 10*time.Second)
- _, err := manager.Discover(discoveryContext)
- cancelDiscovery()
- if err != nil {
- logger.Debug("periodic modem discovery failed", "error", err)
- return
- }
- for _, entry := range manager.List() {
- if !entry.Discovered {
- continue
- }
- refreshContext, cancelRefresh := context.WithTimeout(ctx, 30*time.Second)
- snapshot, err := manager.Refresh(refreshContext, entry.ID)
- cancelRefresh()
- if err != nil && ctx.Err() == nil {
- logger.Warn("modem snapshot refresh failed", "device_id", entry.ID, "error", err)
- }
- if err == nil && ctx.Err() == nil {
- enforceCardRegion(ctx, logger, database, manager, entry.ID, &snapshot)
- }
- }
- }
- refresh()
- ticker := time.NewTicker(30 * time.Second)
- defer ticker.Stop()
- for {
- select {
- case <-ctx.Done():
- return
- case <-ticker.C:
- refresh()
- }
- }
- }
- // cardPolicySourceRegionBlock marks a card policy that was written automatically
- // because the inserted SIM belongs to a region the product does not serve. It
- // doubles as the persistent record that the radio was forced off by us, so the
- // block survives restarts and can be lifted when an allowed card is detected.
- const cardPolicySourceRegionBlock = "auto_region_block"
- // enforceCardRegion applies the regional service policy for one refreshed
- // device. A SIM whose IMSI home MCC is blocked (mainland China, 460/461) is
- // denied service: the radio is forced into airplane mode and a blocking card
- // policy is persisted. The check is fail-open — it only acts on a positively
- // read blocked IMSI — and the lift path only runs once the current card is
- // positively confirmed to be allowed, so an unreadable IMSI never causes a
- // block or a spurious restore.
- func enforceCardRegion(
- ctx context.Context,
- logger *slog.Logger,
- database *store.Store,
- manager *device.Manager,
- id string,
- snapshot *device.Snapshot,
- ) {
- if snapshot == nil || !snapshot.SIMReady {
- return
- }
- imsi := strings.TrimSpace(snapshot.IMSI)
- if imsi == "" {
- // Region unknown: hold the current state rather than block or restore.
- return
- }
- if reason := device.RegionBlockReason(imsi); reason != "" {
- if !snapshot.FlightMode {
- flightContext, cancelFlight := context.WithTimeout(ctx, 30*time.Second)
- _, err := manager.SetFlight(flightContext, id, true)
- cancelFlight()
- if err != nil && ctx.Err() == nil {
- logger.Warn(
- "region block: failed to force airplane mode",
- "device_id", id, "error", err,
- )
- }
- }
- if snapshot.ICCID != "" {
- policy := store.CardPolicy{
- ICCID: snapshot.ICCID,
- NetworkEnabled: false,
- VoWiFiEnabled: false,
- AirplaneEnabled: true,
- IPVersion: "IPV4V6",
- Source: cardPolicySourceRegionBlock,
- }
- if err := database.UpsertCardPolicy(ctx, policy); err != nil && ctx.Err() == nil {
- logger.Warn(
- "region block: failed to persist card policy",
- "device_id", id, "iccid", snapshot.ICCID, "error", err,
- )
- }
- }
- logger.Warn(
- "blocked SIM detected; service disabled and radio forced off",
- "device_id", id, "iccid", snapshot.ICCID, "imsi", imsi, "reason", reason,
- )
- return
- }
- liftCardRegionBlock(ctx, logger, database, manager, id, snapshot)
- }
- // liftCardRegionBlock reverses an automatic region block once the current SIM
- // is positively confirmed to be allowed. It restores the radio only when an
- // outstanding auto-forced block exists, so it never overrides a flight mode the
- // user enabled deliberately.
- func liftCardRegionBlock(
- ctx context.Context,
- logger *slog.Logger,
- database *store.Store,
- manager *device.Manager,
- id string,
- snapshot *device.Snapshot,
- ) {
- policies, err := database.ListCardPolicies(ctx)
- if err != nil {
- if ctx.Err() == nil {
- logger.Warn("region block: failed to list card policies", "error", err)
- }
- return
- }
- outstanding := make([]store.CardPolicy, 0, 1)
- for _, policy := range policies {
- if policy.Source == cardPolicySourceRegionBlock {
- outstanding = append(outstanding, policy)
- }
- }
- if len(outstanding) == 0 {
- return
- }
- if snapshot.FlightMode {
- flightContext, cancelFlight := context.WithTimeout(ctx, 30*time.Second)
- _, err := manager.SetFlight(flightContext, id, false)
- cancelFlight()
- if err != nil && ctx.Err() == nil {
- logger.Warn(
- "region block: failed to restore radio",
- "device_id", id, "error", err,
- )
- return
- }
- }
- for _, policy := range outstanding {
- if err := database.DeleteCardPolicy(ctx, policy.ICCID); err != nil && ctx.Err() == nil {
- logger.Warn(
- "region block: failed to clear auto policy",
- "iccid", policy.ICCID, "error", err,
- )
- }
- }
- logger.Info(
- "region block lifted; SIM is allowed",
- "device_id", id, "iccid", snapshot.ICCID, "imsi", snapshot.IMSI,
- )
- }
|