Capture, decrypt, and restream toolkit with compiled-in app modules (RTE, TG4, BBC), on-device MITM proxy, streamd control plane, and www. BBC module.yaml is published (clear streams); other module values stay local.
118 lines
2.5 KiB
Go
118 lines
2.5 KiB
Go
package device
|
|
|
|
import (
|
|
"fmt"
|
|
"sync"
|
|
|
|
"drmdecryption/adb"
|
|
)
|
|
|
|
// Pool tracks free/busy ADB devices. Never installs apps.
|
|
type Pool struct {
|
|
mu sync.Mutex
|
|
adbBin string
|
|
allow map[string]bool // empty allow = all serials
|
|
busy map[string]string // serial → job label
|
|
}
|
|
|
|
func New(adbBin string, allowlist []string) *Pool {
|
|
p := &Pool{
|
|
adbBin: adbBin,
|
|
allow: map[string]bool{},
|
|
busy: map[string]string{},
|
|
}
|
|
for _, s := range allowlist {
|
|
if s != "" {
|
|
p.allow[s] = true
|
|
}
|
|
}
|
|
return p
|
|
}
|
|
|
|
type Info struct {
|
|
Serial string `json:"serial"`
|
|
State string `json:"state"`
|
|
Model string `json:"model"`
|
|
Busy bool `json:"busy"`
|
|
BusyFor string `json:"busy_for,omitempty"`
|
|
Allowed bool `json:"allowed"`
|
|
}
|
|
|
|
func (p *Pool) List() ([]Info, error) {
|
|
devs, err := adb.ListDevices(p.adbBin)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
p.mu.Lock()
|
|
defer p.mu.Unlock()
|
|
out := make([]Info, 0, len(devs))
|
|
for _, d := range devs {
|
|
allowed := len(p.allow) == 0 || p.allow[d.Serial]
|
|
busyFor, busy := p.busy[d.Serial]
|
|
out = append(out, Info{
|
|
Serial: d.Serial,
|
|
State: d.State,
|
|
Model: d.Model,
|
|
Busy: busy,
|
|
BusyFor: busyFor,
|
|
Allowed: allowed,
|
|
})
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// Acquire finds a free device in state "device" that already has pkg installed.
|
|
// Returns ErrWait if none available (do not burn job attempts).
|
|
func (p *Pool) Acquire(pkg, jobLabel string) (*adb.Client, string, error) {
|
|
devs, err := adb.ListDevices(p.adbBin)
|
|
if err != nil {
|
|
return nil, "", err
|
|
}
|
|
base := adb.New()
|
|
if p.adbBin != "" {
|
|
base.Bin = p.adbBin
|
|
}
|
|
|
|
p.mu.Lock()
|
|
defer p.mu.Unlock()
|
|
|
|
var lastMiss string
|
|
for _, d := range devs {
|
|
if d.State != "device" {
|
|
continue
|
|
}
|
|
if len(p.allow) > 0 && !p.allow[d.Serial] {
|
|
continue
|
|
}
|
|
if _, busy := p.busy[d.Serial]; busy {
|
|
continue
|
|
}
|
|
c := base.WithSerial(d.Serial)
|
|
if pkg != "" && !c.PackageInstalled(pkg) {
|
|
lastMiss = fmt.Sprintf("%s missing package %s", d.Serial, pkg)
|
|
continue
|
|
}
|
|
p.busy[d.Serial] = jobLabel
|
|
return c, d.Serial, nil
|
|
}
|
|
if lastMiss != "" {
|
|
return nil, "", &WaitError{Msg: lastMiss}
|
|
}
|
|
return nil, "", &WaitError{Msg: "no free adb device"}
|
|
}
|
|
|
|
func (p *Pool) Release(serial string) {
|
|
p.mu.Lock()
|
|
defer p.mu.Unlock()
|
|
delete(p.busy, serial)
|
|
}
|
|
|
|
// WaitError means the job should stay queued without consuming an attempt.
|
|
type WaitError struct{ Msg string }
|
|
|
|
func (e *WaitError) Error() string { return e.Msg }
|
|
|
|
func IsWait(err error) bool {
|
|
_, ok := err.(*WaitError)
|
|
return ok
|
|
}
|