Skip to content

Commit 455dc28

Browse files
committed
[client] serve relay control events and account creds over AppControl gRPC
1 parent 4fec13f commit 455dc28

3 files changed

Lines changed: 148 additions & 21 deletions

File tree

client/appcontrol.go

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -117,6 +117,31 @@ func (h *telemetryBroadcaster) unsubscribe(ch chan int32) {
117117
close(ch)
118118
}
119119

120+
func (s *appControlServer) StreamEvents(_ *appcontrolpb.StreamEventsRequest, stream grpc.ServerStreamingServer[appcontrolpb.ProxyEvent]) error {
121+
ch := eventHub.subscribe()
122+
defer eventHub.unsubscribe(ch)
123+
ctx := stream.Context()
124+
for {
125+
select {
126+
case <-ctx.Done():
127+
return nil
128+
case ev, ok := <-ch:
129+
if !ok {
130+
return nil
131+
}
132+
if err := stream.Send(ev); err != nil {
133+
return err
134+
}
135+
}
136+
}
137+
}
138+
139+
func (s *appControlServer) SubmitVKAccountCreds(_ context.Context, req *appcontrolpb.VKAccountCredsRequest) (*appcontrolpb.VKAccountCredsResponse, error) {
140+
log.Printf("app-control: SubmitVKAccountCreds link=%q urls=%d cancel=%t", req.GetLink(), len(req.GetTurnUrls()), req.GetCancel())
141+
dispatchVKAccountCreds(req.GetLink(), req.GetUsername(), req.GetCredential(), req.GetTurnUrls(), req.GetCancel())
142+
return &appcontrolpb.VKAccountCredsResponse{}, nil
143+
}
144+
120145
func (s *appControlServer) GetVKCookies(context.Context, *appcontrolpb.GetVKCookiesRequest) (*appcontrolpb.GetVKCookiesResponse, error) {
121146
cookies, ua, _ := getVkSession()
122147
log.Printf("app-control: GetVKCookies (cookies=%d bytes)", len(cookies))

client/events.go

Lines changed: 94 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,10 @@ import (
55
"fmt"
66
"log"
77
"strings"
8+
"sync"
89
"time"
10+
11+
"github.com/cacggghp/vk-turn-proxy/appcontrolpb"
912
)
1013

1114
const proxyEventProtocolVersion = 1
@@ -91,7 +94,98 @@ func emitProxyEvent(payload any) {
9194
log.Printf("failed to marshal proxy event: %s", err)
9295
return
9396
}
97+
// The JSONL line stays on stdout for the log. Control now flows over the
98+
// AppControl StreamEvents gRPC channel: publish the same event typed there.
9499
fmt.Println("PROXY_EVENT: " + string(encoded))
100+
publishProxyEvent(payload)
101+
}
102+
103+
// proxyEventBroadcaster fans relay control events out to every active
104+
// StreamEvents subscriber. Sends are non-blocking: a subscriber that falls behind
105+
// drops events rather than stalling the relay - control events are re-emitted
106+
// (captcha / vk-auth) or superseded (status), so a rare drop is tolerable.
107+
type proxyEventBroadcaster struct {
108+
mu sync.Mutex
109+
subs map[chan *appcontrolpb.ProxyEvent]struct{}
110+
}
111+
112+
var eventHub = &proxyEventBroadcaster{subs: make(map[chan *appcontrolpb.ProxyEvent]struct{})}
113+
114+
func (h *proxyEventBroadcaster) publish(ev *appcontrolpb.ProxyEvent) {
115+
h.mu.Lock()
116+
defer h.mu.Unlock()
117+
for ch := range h.subs {
118+
select {
119+
case ch <- ev:
120+
default:
121+
}
122+
}
123+
}
124+
125+
func (h *proxyEventBroadcaster) subscribe() chan *appcontrolpb.ProxyEvent {
126+
ch := make(chan *appcontrolpb.ProxyEvent, 64)
127+
h.mu.Lock()
128+
h.subs[ch] = struct{}{}
129+
h.mu.Unlock()
130+
return ch
131+
}
132+
133+
func (h *proxyEventBroadcaster) unsubscribe(ch chan *appcontrolpb.ProxyEvent) {
134+
h.mu.Lock()
135+
delete(h.subs, ch)
136+
h.mu.Unlock()
137+
close(ch)
138+
}
139+
140+
// publishProxyEvent maps a JSONL event struct to its typed ProxyEvent and pushes
141+
// it to StreamEvents subscribers. Payloads with no control mapping (telemetry,
142+
// vk_cookies_update) are logged as JSONL only and skipped here - telemetry has its
143+
// own StreamTelemetry RPC and the app pulls rotated cookies via GetVKCookies.
144+
func publishProxyEvent(payload any) {
145+
var ev *appcontrolpb.ProxyEvent
146+
switch p := payload.(type) {
147+
case proxyStatusEvent:
148+
ev = &appcontrolpb.ProxyEvent{Event: &appcontrolpb.ProxyEvent_Status{
149+
Status: &appcontrolpb.StatusEvent{Phase: p.Phase},
150+
}}
151+
case proxyStreamStatusEvent:
152+
ev = &appcontrolpb.ProxyEvent{Event: &appcontrolpb.ProxyEvent_Status{
153+
Status: &appcontrolpb.StatusEvent{Phase: p.Phase, StreamId: int32(p.StreamID)},
154+
}}
155+
case proxyDtlsAliveEvent:
156+
ev = &appcontrolpb.ProxyEvent{Event: &appcontrolpb.ProxyEvent_Status{
157+
Status: &appcontrolpb.StatusEvent{Phase: p.Phase, StreamId: int32(p.StreamID), ConnectedStreams: int32(p.ConnectedStreams)},
158+
}}
159+
case proxyCapsEvent:
160+
ev = &appcontrolpb.ProxyEvent{Event: &appcontrolpb.ProxyEvent_Caps{
161+
Caps: &appcontrolpb.CapsEvent{Version: int32(p.Version), Capabilities: append([]string(nil), p.Capabilities...)},
162+
}}
163+
case proxyCaptchaEvent:
164+
ev = &appcontrolpb.ProxyEvent{Event: &appcontrolpb.ProxyEvent_Captcha{
165+
Captcha: &appcontrolpb.CaptchaEvent{State: p.State, Source: p.Source, Url: p.URL, UserAgent: p.UserAgent},
166+
}}
167+
case proxyLockoutEvent:
168+
ev = &appcontrolpb.ProxyEvent{Event: &appcontrolpb.ProxyEvent_Lockout{
169+
Lockout: &appcontrolpb.LockoutEvent{Seconds: int32(p.Seconds)},
170+
}}
171+
case vkAccountAuthEvent:
172+
if p.Type == "vk_cookies_required" {
173+
ev = &appcontrolpb.ProxyEvent{Event: &appcontrolpb.ProxyEvent_VkCookiesRequired{
174+
VkCookiesRequired: &appcontrolpb.VKCookiesRequiredEvent{},
175+
}}
176+
} else {
177+
ev = &appcontrolpb.ProxyEvent{Event: &appcontrolpb.ProxyEvent_VkAccountAuth{
178+
VkAccountAuth: &appcontrolpb.VKAccountAuthEvent{
179+
Phase: strings.TrimPrefix(p.Type, "vk_account_auth_"),
180+
Link: p.Link,
181+
Reason: p.Reason,
182+
},
183+
}}
184+
}
185+
default:
186+
return
187+
}
188+
eventHub.publish(ev)
95189
}
96190

97191
func applyProxyStatusState(marker string) {

client/vk_account.go

Lines changed: 29 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -320,28 +320,36 @@ func StartAccountCredsStdinReader(ctx context.Context) {
320320
if msg.Type != "vk_account_creds" {
321321
continue
322322
}
323-
link := vkHashFromLink(msg.Link)
324-
if link == "" {
325-
continue
326-
}
327-
if msg.Cancel {
328-
log.Printf("[VK Auth] account creds cancel received for link %s", link)
329-
resolveAccountAuth(link, accountCredsResult{err: fmt.Errorf("VK account auth cancelled by app")})
330-
continue
331-
}
332-
addresses := turnURLsToAddresses(msg.URLs)
333-
if msg.Username == "" || msg.Credential == "" || len(addresses) == 0 {
334-
log.Printf("[VK Auth] ignoring incomplete account creds line for link %s", link)
335-
continue
336-
}
337-
injectTurnCreds(link, msg.Username, msg.Credential, msg.URLs)
338-
resolveAccountAuth(link, accountCredsResult{creds: injectedTurnCreds{
339-
user: msg.Username,
340-
pass: msg.Credential,
341-
addrs: cloneAddrs(addresses),
342-
}})
343-
log.Printf("[VK Auth] received account TURN creds for link %s (urls=%d)", link, len(addresses))
323+
dispatchVKAccountCreds(msg.Link, msg.Username, msg.Credential, msg.URLs, msg.Cancel)
344324
}
345325
log.Printf("[VK Auth] account creds stdin reader exited (err=%v)", scanner.Err())
346326
}()
347327
}
328+
329+
// dispatchVKAccountCreds resolves a link's pending account-auth wait with the VK
330+
// TURN creds the host app intercepted (or a cancel). Shared by the stdin reader
331+
// and the AppControl SubmitVKAccountCreds RPC so both delivery paths behave
332+
// identically.
333+
func dispatchVKAccountCreds(rawLink, username, credential string, urls []string, cancel bool) {
334+
link := vkHashFromLink(rawLink)
335+
if link == "" {
336+
return
337+
}
338+
if cancel {
339+
log.Printf("[VK Auth] account creds cancel received for link %s", link)
340+
resolveAccountAuth(link, accountCredsResult{err: fmt.Errorf("VK account auth cancelled by app")})
341+
return
342+
}
343+
addresses := turnURLsToAddresses(urls)
344+
if username == "" || credential == "" || len(addresses) == 0 {
345+
log.Printf("[VK Auth] ignoring incomplete account creds for link %s", link)
346+
return
347+
}
348+
injectTurnCreds(link, username, credential, urls)
349+
resolveAccountAuth(link, accountCredsResult{creds: injectedTurnCreds{
350+
user: username,
351+
pass: credential,
352+
addrs: cloneAddrs(addresses),
353+
}})
354+
log.Printf("[VK Auth] received account TURN creds for link %s (urls=%d)", link, len(addresses))
355+
}

0 commit comments

Comments
 (0)