main.go 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617
  1. package main
  2. import (
  3. "context"
  4. "encoding/json"
  5. "errors"
  6. "fmt"
  7. "log/slog"
  8. "net/http"
  9. "os"
  10. "os/signal"
  11. "strings"
  12. "syscall"
  13. "time"
  14. "vocat/internal/auth"
  15. "vocat/internal/config"
  16. "vocat/internal/device"
  17. "vocat/internal/loghub"
  18. "vocat/internal/server"
  19. "vocat/internal/store"
  20. "vocat/internal/update"
  21. "vocat/internal/vowifi"
  22. "vocat/internal/vowifi/ike"
  23. "vocat/internal/vowifi/ims"
  24. "vocat/internal/vowifi/integration"
  25. vowifiruntime "vocat/internal/vowifi/runtime"
  26. "vocat/web"
  27. )
  28. func main() {
  29. logs := loghub.New(slog.NewJSONHandler(os.Stdout, nil), 2000)
  30. logger := slog.New(logs)
  31. args := os.Args[1:]
  32. switch subcommand, rest := splitSubcommand(args); subcommand {
  33. case "":
  34. // No subcommand: run the server. Backward-compatible with the
  35. // existing systemd unit (ExecStart=/opt/vocat/bin/vocat).
  36. if err := run(logger, logs); err != nil {
  37. logger.Error("server stopped", "error", err)
  38. os.Exit(1)
  39. }
  40. case "version", "-v", "--version":
  41. runVersion()
  42. case "update":
  43. if err := update.Run(logger, rest); err != nil {
  44. logger.Error("update failed", "error", err)
  45. os.Exit(1)
  46. }
  47. case "menu":
  48. if err := runMenu(logger); err != nil {
  49. logger.Error("menu failed", "error", err)
  50. os.Exit(1)
  51. }
  52. case "help", "-h", "--help":
  53. printUsage(os.Stdout)
  54. default:
  55. fmt.Fprintf(os.Stderr, "vocat: unknown subcommand %q\n\n", subcommand)
  56. printUsage(os.Stderr)
  57. os.Exit(2)
  58. }
  59. }
  60. // splitSubcommand returns the first non-flag token as the subcommand and the
  61. // remaining args. An empty arg list yields ("", nil) → server mode.
  62. func splitSubcommand(args []string) (string, []string) {
  63. if len(args) == 0 {
  64. return "", nil
  65. }
  66. return args[0], args[1:]
  67. }
  68. func run(logger *slog.Logger, logs *loghub.Hub) error {
  69. cfg, err := config.Load()
  70. if err != nil {
  71. return fmt.Errorf("load configuration: %w", err)
  72. }
  73. if cfg.UsesDefaultCredentials() {
  74. logger.Warn(
  75. "default admin credentials are active; set VOCAT_ADMIN_PASSWORD before exposing the service",
  76. )
  77. }
  78. startupContext, cancelStartup := context.WithTimeout(context.Background(), 15*time.Second)
  79. defer cancelStartup()
  80. database, err := store.Open(startupContext, cfg.DatabasePath)
  81. if err != nil {
  82. return err
  83. }
  84. defer database.Close()
  85. authService, err := auth.New(database, auth.Options{
  86. SessionTTL: cfg.SessionTTL,
  87. })
  88. if err != nil {
  89. return err
  90. }
  91. if err := authService.EnsureAdmin(
  92. startupContext,
  93. cfg.AdminUsername,
  94. cfg.AdminPassword,
  95. ); err != nil {
  96. return err
  97. }
  98. deviceManager, err := device.NewManager(device.Options{})
  99. if err != nil {
  100. return fmt.Errorf("create device manager: %w", err)
  101. }
  102. if err := deviceManager.Start(startupContext); err != nil {
  103. logger.Warn("device discovery is not available at startup", "error", err)
  104. }
  105. if err := provisionDiscoveredDevices(startupContext, database, deviceManager); err != nil {
  106. logger.Warn("automatic first-run device provisioning failed", "error", err)
  107. }
  108. defer func() {
  109. stopContext, cancel := context.WithTimeout(context.Background(), 5*time.Second)
  110. defer cancel()
  111. if err := deviceManager.Stop(stopContext); err != nil {
  112. logger.Warn("stop device manager", "error", err)
  113. }
  114. }()
  115. pollContext, cancelPolling := context.WithCancel(context.Background())
  116. defer cancelPolling()
  117. go pollDeviceSnapshots(pollContext, logger, database, deviceManager)
  118. go persistLogsToStore(pollContext, logger, logs, database)
  119. vowifiManager, err := configureVoWiFiRuntime(
  120. startupContext,
  121. logger,
  122. database,
  123. deviceManager,
  124. )
  125. if err != nil {
  126. return fmt.Errorf("configure VoWiFi runtime: %w", err)
  127. }
  128. defer func() {
  129. stopContext, cancel := context.WithTimeout(context.Background(), 15*time.Second)
  130. defer cancel()
  131. if err := vowifiManager.Close(stopContext); err != nil {
  132. logger.Warn("stop VoWiFi runtime", "error", err)
  133. }
  134. }()
  135. handler, err := server.New(server.Options{
  136. Store: database,
  137. Auth: authService,
  138. Devices: deviceManager,
  139. VoWiFi: vowifiManager,
  140. Logs: logs,
  141. Assets: web.Dist,
  142. Logger: logger,
  143. SecureCookies: cfg.SecureCookies,
  144. MaxRequestBodyBytes: cfg.MaxRequestBodyBytes,
  145. })
  146. if err != nil {
  147. return err
  148. }
  149. go handler.StartLogRetentionLoop(pollContext, time.Minute)
  150. go handler.StartSMSSyncLoop(pollContext, 15*time.Second)
  151. handler.StartTelegramBot(pollContext)
  152. handler.StartSMSNotificationDispatchers(pollContext)
  153. httpServer := &http.Server{
  154. Addr: cfg.Address,
  155. Handler: handler,
  156. ReadHeaderTimeout: 5 * time.Second,
  157. ReadTimeout: 15 * time.Second,
  158. WriteTimeout: 30 * time.Second,
  159. IdleTimeout: 90 * time.Second,
  160. MaxHeaderBytes: 1 << 20,
  161. }
  162. signalContext, stopSignals := signal.NotifyContext(
  163. context.Background(),
  164. os.Interrupt,
  165. syscall.SIGTERM,
  166. )
  167. defer stopSignals()
  168. serverError := make(chan error, 1)
  169. go func() {
  170. logger.Info("HTTP server listening", "address", cfg.Address)
  171. err := httpServer.ListenAndServe()
  172. if errors.Is(err, http.ErrServerClosed) {
  173. err = nil
  174. }
  175. serverError <- err
  176. }()
  177. select {
  178. case err := <-serverError:
  179. return err
  180. case <-signalContext.Done():
  181. logger.Info("shutdown signal received")
  182. }
  183. shutdownContext, cancelShutdown := context.WithTimeout(
  184. context.Background(),
  185. cfg.ShutdownTimeout,
  186. )
  187. defer cancelShutdown()
  188. if err := httpServer.Shutdown(shutdownContext); err != nil {
  189. _ = httpServer.Close()
  190. return fmt.Errorf("graceful HTTP shutdown: %w", err)
  191. }
  192. return <-serverError
  193. }
  194. func configureVoWiFiRuntime(
  195. ctx context.Context,
  196. logger *slog.Logger,
  197. database *store.Store,
  198. deviceManager *device.Manager,
  199. ) (*vowifiruntime.Manager, error) {
  200. mapper := integration.ATMapper{
  201. Store: database,
  202. Devices: deviceManager,
  203. }
  204. adapter, err := vowifi.NewEC20Adapter(mapper, vowifi.EC20AdapterOptions{
  205. // The test deployment is deliberately non-cellular. VoWiFi teardown
  206. // may restore CFUN, but it must never reactivate a PDP context.
  207. RestoreCellularData: false,
  208. })
  209. if err != nil {
  210. return nil, err
  211. }
  212. projector := integration.StateProjector{
  213. Store: database,
  214. Devices: mapper,
  215. }
  216. manager := vowifiruntime.New(vowifiruntime.Options{
  217. Logger: logger,
  218. OnState: projector.Save,
  219. Factory: func(factoryContext context.Context, deviceID string) (*vowifi.Orchestrator, error) {
  220. deviceConfig, err := database.Device(factoryContext, deviceID)
  221. if err != nil {
  222. return nil, fmt.Errorf("load device %q VoWiFi config: %w", deviceID, err)
  223. }
  224. return newVoWiFiOrchestrator(deviceConfig, database, adapter)
  225. },
  226. })
  227. configured, err := database.ListDevices(ctx)
  228. if err != nil {
  229. _ = manager.Close(context.Background())
  230. return nil, err
  231. }
  232. for _, deviceConfig := range configured {
  233. if err := manager.Ensure(ctx, deviceConfig.ID); err != nil {
  234. _ = manager.Close(context.Background())
  235. return nil, fmt.Errorf("register device %q VoWiFi runtime: %w", deviceConfig.ID, err)
  236. }
  237. if deviceConfig.VoWiFiEnabled {
  238. if _, err := manager.RequestEnabled(deviceConfig.ID, true); err != nil {
  239. _ = manager.Close(context.Background())
  240. return nil, fmt.Errorf("start device %q VoWiFi policy: %w", deviceConfig.ID, err)
  241. }
  242. }
  243. }
  244. return manager, nil
  245. }
  246. func newVoWiFiOrchestrator(
  247. deviceConfig store.Device,
  248. database *store.Store,
  249. adapter *vowifi.EC20Adapter,
  250. ) (*vowifi.Orchestrator, error) {
  251. apn := deviceConfig.APN
  252. if apn == "" {
  253. apn = "ims"
  254. }
  255. tunnelProvider, err := ike.NewProvider(ike.Config{APN: apn})
  256. if err != nil {
  257. return nil, fmt.Errorf("device %q IKE provider: %w", deviceConfig.ID, err)
  258. }
  259. imsProvider, err := ims.NewProvider(adapter, ims.Config{
  260. // The userspace SWu data plane currently carries the protected P-CSCF
  261. // signalling path over TCP.
  262. Transport: "tcp",
  263. // Some Vodafone UK SIM profiles leave AT+CSCA empty; Vodafone publishes
  264. // this service-centre number for manual SMS setup.
  265. SMSCenter: "+447785016005",
  266. OnSMS: func(ctx context.Context, message ims.ReceivedSMS) error {
  267. extra, _ := json.Marshal(map[string]any{
  268. "transport": "ims",
  269. "encoding": message.Encoding,
  270. "concat": message.Concat,
  271. "rp_reference": message.RPReference,
  272. "call_id": message.CallID,
  273. "received_at": message.Timestamp,
  274. "service_center_timestamp": message.ServiceCenterTimestamp,
  275. "raw_rpdu": message.RawRPDU,
  276. "raw_tpdu": message.RawTPDU,
  277. })
  278. partsTotal := 1
  279. if message.Concat != nil && message.Concat.Total > 0 {
  280. partsTotal = message.Concat.Total
  281. }
  282. _, saveErr := database.SaveSMSMessage(ctx, store.SMSMessage{
  283. MessageID: message.MessageID,
  284. DeviceID: message.DeviceID,
  285. IMSI: message.IMSI,
  286. Peer: message.From,
  287. Direction: "inbound",
  288. Body: message.Text,
  289. Timestamp: message.Timestamp,
  290. Status: "received",
  291. Source: "ims",
  292. PartsTotal: partsTotal,
  293. Read: false,
  294. Extra: extra,
  295. })
  296. return saveErr
  297. },
  298. OnSMSStatus: func(ctx context.Context, report ims.ReceivedSMSStatus) error {
  299. deliveryReport := store.SMSDeliveryReport{
  300. DeviceID: report.DeviceID,
  301. IMSI: report.IMSI,
  302. Peer: report.To,
  303. Source: "ims",
  304. MessageReference: report.MessageReference,
  305. StatusCode: report.StatusCode,
  306. DeliveryState: report.DeliveryStatus,
  307. ServiceCenterTime: report.ServiceCenterTimestamp,
  308. DischargeTime: report.DischargeTimestamp,
  309. ReceivedAt: report.Timestamp,
  310. }
  311. var applyErr error
  312. for attempt := 0; attempt < 10; attempt++ {
  313. _, applyErr = database.ApplySMSDeliveryReport(ctx, deliveryReport)
  314. if !errors.Is(applyErr, store.ErrNotFound) {
  315. return applyErr
  316. }
  317. // A status report can race the API handler persisting the SIP 202
  318. // result. Give that write a brief chance to complete.
  319. select {
  320. case <-ctx.Done():
  321. return ctx.Err()
  322. case <-time.After(100 * time.Millisecond):
  323. }
  324. }
  325. // A late report from before this process started must still be
  326. // acknowledged, otherwise the SMSC will keep retransmitting it.
  327. return nil
  328. },
  329. })
  330. if err != nil {
  331. return nil, fmt.Errorf("device %q IMS provider: %w", deviceConfig.ID, err)
  332. }
  333. orchestrator, err := vowifi.New(vowifi.Dependencies{
  334. SIM: adapter,
  335. AKA: adapter,
  336. Radio: adapter,
  337. Proxy: integration.ProxyResolver{Store: database},
  338. Tunnel: tunnelProvider,
  339. IMS: imsProvider,
  340. Phones: integration.PhoneStore{Store: database, DeviceID: deviceConfig.ID},
  341. }, vowifi.Options{
  342. DeviceID: deviceConfig.ID,
  343. AllowIMSWithoutSMS: true,
  344. })
  345. if err != nil {
  346. return nil, fmt.Errorf("device %q VoWiFi orchestrator: %w", deviceConfig.ID, err)
  347. }
  348. return orchestrator, nil
  349. }
  350. func provisionDiscoveredDevices(
  351. ctx context.Context,
  352. database *store.Store,
  353. manager *device.Manager,
  354. ) error {
  355. configured, err := database.ListDevices(ctx)
  356. if err != nil {
  357. return err
  358. }
  359. if len(configured) != 0 {
  360. return nil
  361. }
  362. for _, discovered := range manager.List() {
  363. candidate := discovered.Candidate
  364. backend := "at"
  365. control := candidate.ATPort.OpenPath()
  366. if candidate.QMIControl != "" {
  367. backend = "qmi"
  368. control = candidate.QMIControl
  369. }
  370. name := candidate.Product
  371. if name == "" || strings.EqualFold(name, "Android") {
  372. name = "Quectel EC20 / EC25"
  373. }
  374. if err := database.UpsertDevice(ctx, store.Device{
  375. ID: discovered.ID,
  376. Name: name,
  377. Interface: candidate.NetworkInterface,
  378. ControlDevice: control,
  379. ATPort: candidate.ATPort.OpenPath(),
  380. USBPath: candidate.USBPath,
  381. ProxyPort: 1080,
  382. BaudRate: 115200,
  383. DataBits: 8,
  384. StopBits: 1,
  385. Parity: "none",
  386. DeviceBackend: backend,
  387. ESIMTransport: backend,
  388. NetworkEnabled: false,
  389. SMSEnabled: true,
  390. VoWiFiEnabled: false,
  391. }); err != nil {
  392. return err
  393. }
  394. }
  395. return nil
  396. }
  397. // persistLogsToStore subscribes to the live log hub and durably appends every
  398. // entry to the log_events table, so runtime logs survive restarts and can be
  399. // pruned by the configured retention policy.
  400. func persistLogsToStore(
  401. ctx context.Context,
  402. logger *slog.Logger,
  403. logs *loghub.Hub,
  404. database *store.Store,
  405. ) {
  406. entries, cancel := logs.Subscribe(512)
  407. defer cancel()
  408. for {
  409. select {
  410. case <-ctx.Done():
  411. return
  412. case entry, ok := <-entries:
  413. if !ok {
  414. return
  415. }
  416. var fields json.RawMessage
  417. if len(entry.Fields) > 0 {
  418. if raw, err := json.Marshal(entry.Fields); err == nil {
  419. fields = raw
  420. }
  421. }
  422. if _, err := database.AppendLogEvent(ctx, store.LogEvent{
  423. Time: entry.Time,
  424. Level: entry.Level,
  425. Message: entry.Message,
  426. Caller: entry.Caller,
  427. Fields: fields,
  428. }); err != nil && ctx.Err() == nil {
  429. logger.Warn("persist log event failed", "error", err)
  430. }
  431. }
  432. }
  433. }
  434. func pollDeviceSnapshots(
  435. ctx context.Context,
  436. logger *slog.Logger,
  437. database *store.Store,
  438. manager *device.Manager,
  439. ) {
  440. refresh := func() {
  441. discoveryContext, cancelDiscovery := context.WithTimeout(ctx, 10*time.Second)
  442. _, err := manager.Discover(discoveryContext)
  443. cancelDiscovery()
  444. if err != nil {
  445. logger.Debug("periodic modem discovery failed", "error", err)
  446. return
  447. }
  448. for _, entry := range manager.List() {
  449. if !entry.Discovered {
  450. continue
  451. }
  452. refreshContext, cancelRefresh := context.WithTimeout(ctx, 30*time.Second)
  453. snapshot, err := manager.Refresh(refreshContext, entry.ID)
  454. cancelRefresh()
  455. if err != nil && ctx.Err() == nil {
  456. logger.Warn("modem snapshot refresh failed", "device_id", entry.ID, "error", err)
  457. }
  458. if err == nil && ctx.Err() == nil {
  459. enforceCardRegion(ctx, logger, database, manager, entry.ID, &snapshot)
  460. }
  461. }
  462. }
  463. refresh()
  464. ticker := time.NewTicker(30 * time.Second)
  465. defer ticker.Stop()
  466. for {
  467. select {
  468. case <-ctx.Done():
  469. return
  470. case <-ticker.C:
  471. refresh()
  472. }
  473. }
  474. }
  475. // cardPolicySourceRegionBlock marks a card policy that was written automatically
  476. // because the inserted SIM belongs to a region the product does not serve. It
  477. // doubles as the persistent record that the radio was forced off by us, so the
  478. // block survives restarts and can be lifted when an allowed card is detected.
  479. const cardPolicySourceRegionBlock = "auto_region_block"
  480. // enforceCardRegion applies the regional service policy for one refreshed
  481. // device. A SIM whose IMSI home MCC is blocked (mainland China, 460/461) is
  482. // denied service: the radio is forced into airplane mode and a blocking card
  483. // policy is persisted. The check is fail-open — it only acts on a positively
  484. // read blocked IMSI — and the lift path only runs once the current card is
  485. // positively confirmed to be allowed, so an unreadable IMSI never causes a
  486. // block or a spurious restore.
  487. func enforceCardRegion(
  488. ctx context.Context,
  489. logger *slog.Logger,
  490. database *store.Store,
  491. manager *device.Manager,
  492. id string,
  493. snapshot *device.Snapshot,
  494. ) {
  495. if snapshot == nil || !snapshot.SIMReady {
  496. return
  497. }
  498. imsi := strings.TrimSpace(snapshot.IMSI)
  499. if imsi == "" {
  500. // Region unknown: hold the current state rather than block or restore.
  501. return
  502. }
  503. if reason := device.RegionBlockReason(imsi); reason != "" {
  504. if !snapshot.FlightMode {
  505. flightContext, cancelFlight := context.WithTimeout(ctx, 30*time.Second)
  506. _, err := manager.SetFlight(flightContext, id, true)
  507. cancelFlight()
  508. if err != nil && ctx.Err() == nil {
  509. logger.Warn(
  510. "region block: failed to force airplane mode",
  511. "device_id", id, "error", err,
  512. )
  513. }
  514. }
  515. if snapshot.ICCID != "" {
  516. policy := store.CardPolicy{
  517. ICCID: snapshot.ICCID,
  518. NetworkEnabled: false,
  519. VoWiFiEnabled: false,
  520. AirplaneEnabled: true,
  521. IPVersion: "IPV4V6",
  522. Source: cardPolicySourceRegionBlock,
  523. }
  524. if err := database.UpsertCardPolicy(ctx, policy); err != nil && ctx.Err() == nil {
  525. logger.Warn(
  526. "region block: failed to persist card policy",
  527. "device_id", id, "iccid", snapshot.ICCID, "error", err,
  528. )
  529. }
  530. }
  531. logger.Warn(
  532. "blocked SIM detected; service disabled and radio forced off",
  533. "device_id", id, "iccid", snapshot.ICCID, "imsi", imsi, "reason", reason,
  534. )
  535. return
  536. }
  537. liftCardRegionBlock(ctx, logger, database, manager, id, snapshot)
  538. }
  539. // liftCardRegionBlock reverses an automatic region block once the current SIM
  540. // is positively confirmed to be allowed. It restores the radio only when an
  541. // outstanding auto-forced block exists, so it never overrides a flight mode the
  542. // user enabled deliberately.
  543. func liftCardRegionBlock(
  544. ctx context.Context,
  545. logger *slog.Logger,
  546. database *store.Store,
  547. manager *device.Manager,
  548. id string,
  549. snapshot *device.Snapshot,
  550. ) {
  551. policies, err := database.ListCardPolicies(ctx)
  552. if err != nil {
  553. if ctx.Err() == nil {
  554. logger.Warn("region block: failed to list card policies", "error", err)
  555. }
  556. return
  557. }
  558. outstanding := make([]store.CardPolicy, 0, 1)
  559. for _, policy := range policies {
  560. if policy.Source == cardPolicySourceRegionBlock {
  561. outstanding = append(outstanding, policy)
  562. }
  563. }
  564. if len(outstanding) == 0 {
  565. return
  566. }
  567. if snapshot.FlightMode {
  568. flightContext, cancelFlight := context.WithTimeout(ctx, 30*time.Second)
  569. _, err := manager.SetFlight(flightContext, id, false)
  570. cancelFlight()
  571. if err != nil && ctx.Err() == nil {
  572. logger.Warn(
  573. "region block: failed to restore radio",
  574. "device_id", id, "error", err,
  575. )
  576. return
  577. }
  578. }
  579. for _, policy := range outstanding {
  580. if err := database.DeleteCardPolicy(ctx, policy.ICCID); err != nil && ctx.Err() == nil {
  581. logger.Warn(
  582. "region block: failed to clear auto policy",
  583. "iccid", policy.ICCID, "error", err,
  584. )
  585. }
  586. }
  587. logger.Info(
  588. "region block lifted; SIM is allowed",
  589. "device_id", id, "iccid", snapshot.ICCID, "imsi", snapshot.IMSI,
  590. )
  591. }