Null-DRM-Official/apps/agent/internal/device/pool.go
404errordeveloper 2fa8f2435f Initial commit: Null DRM Official
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.
2026-10-06 00:25:35 +02:00

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
}