Repository navigation
Expand file tree
/
Copy pathclient.go
More file actions
286 lines (247 loc) · 8.3 KB
/
Copy pathclient.go
File metadata and controls
286 lines (247 loc) · 8.3 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
package telegram
import (
"context"
"io"
"sync"
"time"
"github.com/cenkalti/backoff/v4"
"go.opentelemetry.io/otel/trace"
"go.uber.org/atomic"
"golang.org/x/sync/singleflight"
"github.com/gotd/log"
"github.com/gotd/td/bin"
"github.com/gotd/td/clock"
"github.com/gotd/td/mtproto"
"github.com/gotd/td/oteltg"
"github.com/gotd/td/pool"
"github.com/gotd/td/session"
"github.com/gotd/td/tdsync"
"github.com/gotd/td/telegram/dcs"
"github.com/gotd/td/telegram/internal/manager"
"github.com/gotd/td/telegram/internal/version"
"github.com/gotd/td/tg"
)
// UpdateHandler will be called on received updates from Telegram.
type UpdateHandler interface {
Handle(ctx context.Context, u tg.UpdatesClass) error
}
// UpdateHandlerFunc type is an adapter to allow the use of
// ordinary function as update handler.
//
// UpdateHandlerFunc(f) is an UpdateHandler that calls f.
type UpdateHandlerFunc func(ctx context.Context, u tg.UpdatesClass) error
// Handle calls f(ctx, u)
func (f UpdateHandlerFunc) Handle(ctx context.Context, u tg.UpdatesClass) error {
return f(ctx, u)
}
type clientStorage interface {
Load(ctx context.Context) (*session.Data, error)
Save(ctx context.Context, data *session.Data) error
}
type clientConn interface {
Run(ctx context.Context) error
Invoke(ctx context.Context, input bin.Encoder, output bin.Decoder) error
Ping(ctx context.Context) error
}
// Client represents a MTProto client to Telegram.
type Client struct {
// Put migration in the header of the structure to ensure 64-bit alignment,
// otherwise it will cause the atomic operation of connsCounter to panic.
// DO NOT change the order of members arbitrarily.
// Ref: https://pkg.go.dev/sync/atomic#pkg-note-BUG
// Connection factory fields.
connsCounter atomic.Int64
create connConstructor // immutable
resolver dcs.Resolver // immutable
onDead func(error) // immutable
newConnBackoff func() backoff.BackOff // immutable
defaultMode manager.ConnMode // immutable
// onConnectionState is called on primary connection state change.
onConnectionState func(ConnectionState) // immutable
// Migration state.
migrationTimeout time.Duration // immutable
migration chan struct{}
// tg provides RPC calls via Client. Uses invoker below.
tg *tg.Client // immutable
// invoker implements tg.Invoker on top of Client and mw.
invoker tg.Invoker // immutable
// mw is list of middlewares used in invoker, can be blank.
mw []Middleware // immutable
// Telegram device information.
device DeviceConfig // immutable
// Schema layer requested via invokeWithLayer.
layer int // immutable
// MTProto options.
opts mtproto.Options // immutable
// DCList state.
// Domain list (for websocket)
domains map[int]string // immutable
// Denotes to use Test DCs.
testDC bool // immutable
// Connection state. Guarded by connMux.
session *pool.SyncSession
cfg *manager.AtomicConfig
conn clientConn
// connChanged is closed and re-created on every primary connection
// replacement, letting invokeConn wait for reconnect and retry.
connChanged chan struct{}
connBackoff atomic.Pointer[backoff.BackOff]
connMux sync.Mutex
// Restart signal channel.
restart chan struct{} // immutable
// Connections to non-primary DC.
subConns map[int]CloseInvoker
subConnsMux sync.Mutex
// Shared CDN pools and handle references.
cdnPools cdnPoolManager
// sessions stores regular non-primary DC sessions.
sessions map[int]*pool.SyncSession
// cdnSessions stores session state for CDN pools separately from regular DCs.
cdnSessions map[int]*pool.SyncSession
sessionsMux sync.Mutex
// CDN public keys loaded from help.getCdnConfig and cached per CDN DC.
cdnKeys []PublicKey
cdnKeysByDC map[int][]PublicKey
cdnKeysSet bool
// cdnKeysGen increments on cache invalidation to avoid storing stale
// singleflight result after fingerprint miss.
cdnKeysGen uint64
cdnKeysMux sync.Mutex
cdnKeysLoad singleflight.Group
// Wrappers for external world, like logs or PRNG.
rand io.Reader // immutable
log log.Helper // immutable
clock clock.Clock // immutable
// Client context. Will be canceled by Run on exit.
ctx context.Context
cancel context.CancelFunc
// Client config.
appID int // immutable
appHash string // immutable
// allowCDN is the explicit downloader policy copied from Options.AllowCDN.
allowCDN bool // immutable
// Session storage.
storage clientStorage // immutable, nillable
// Ready signal channel, sends signal when client connection is ready.
// Resets on reconnect.
ready *tdsync.ResetReady // immutable
// Telegram updates handler.
updateHandler UpdateHandler // immutable
// Denotes that no update mode is enabled.
noUpdatesMode bool // immutable
// Tracing.
tracer trace.Tracer
// onTransfer is called in transfer.
onTransfer AuthTransferHandler
// onSelfError is called on error calling Self().
onSelfError func(ctx context.Context, err error) error
// onSelfSuccess is called on success calling Self().
onSelfSuccess func(self *tg.User)
}
// NewClient creates new unstarted client.
func NewClient(appID int, appHash string, opt Options) *Client {
opt.setDefaults()
mode := manager.ConnModeUpdates
if opt.NoUpdates {
mode = manager.ConnModeData
}
client := &Client{
rand: opt.Random,
log: log.For(opt.Logger),
appID: appID,
appHash: appHash,
allowCDN: opt.AllowCDN,
updateHandler: opt.UpdateHandler,
session: pool.NewSyncSession(pool.Session{
DC: opt.DC,
}),
domains: opt.DCList.Domains,
testDC: opt.DCList.Test,
cfg: manager.NewAtomicConfig(tg.Config{
DCOptions: opt.DCList.Options,
}),
create: defaultConstructor(),
resolver: opt.Resolver,
defaultMode: mode,
newConnBackoff: opt.ReconnectionBackoff,
onDead: opt.OnDead,
onConnectionState: opt.OnConnectionState,
clock: opt.Clock,
device: opt.Device,
layer: opt.Layer,
migrationTimeout: opt.MigrationTimeout,
noUpdatesMode: opt.NoUpdates,
mw: opt.Middlewares,
onTransfer: opt.OnTransfer,
onSelfError: opt.OnSelfError,
onSelfSuccess: opt.OnSelfSuccess,
}
if opt.TracerProvider != nil {
client.tracer = opt.TracerProvider.Tracer(oteltg.Name)
}
client.init()
// Including version into client logger to help with debugging.
if v := version.GetVersion(); v != "" {
client.log = client.log.With(log.String("v", v))
}
if opt.SessionStorage != nil {
client.storage = &session.Loader{
Storage: opt.SessionStorage,
}
}
client.opts = mtproto.Options{
PublicKeys: opt.PublicKeys,
Random: opt.Random,
Logger: opt.Logger,
AckBatchSize: opt.AckBatchSize,
AckInterval: opt.AckInterval,
RetryInterval: opt.RetryInterval,
MaxRetries: opt.MaxRetries,
CompressThreshold: opt.CompressThreshold,
MessageID: opt.MessageID,
ExchangeTimeout: opt.ExchangeTimeout,
DialTimeout: opt.DialTimeout,
// Forward PFS toggles into low-level mtproto connection.
EnablePFS: opt.EnablePFS,
TempKeyTTL: opt.TempKeyTTL,
Clock: opt.Clock,
Types: getTypesMapping(),
Tracer: client.tracer,
}
client.conn = client.createPrimaryConn(nil)
return client
}
// replaceConn replaces primary connection and notifies invokers waiting for
// reconnection (see invokeConn).
//
// Caller must hold connMux.
func (c *Client) replaceConn(conn clientConn) {
c.conn = conn
close(c.connChanged)
c.connChanged = make(chan struct{})
}
// init sets fields which needs explicit initialization, like maps or channels.
func (c *Client) init() {
if c.domains == nil {
c.domains = map[int]string{}
}
if c.cfg == nil {
c.cfg = manager.NewAtomicConfig(tg.Config{})
}
c.ready = tdsync.NewResetReady()
c.connChanged = make(chan struct{})
c.restart = make(chan struct{})
c.migration = make(chan struct{}, 1)
c.sessions = map[int]*pool.SyncSession{}
c.cdnSessions = map[int]*pool.SyncSession{}
c.subConns = map[int]CloseInvoker{}
c.cdnPools = newCDNPoolManager()
// CDN key cache is cold-started and filled lazily on first CDN pool create.
c.cdnKeys = nil
c.cdnKeysByDC = nil
c.cdnKeysSet = false
c.cdnKeysGen = 0
c.cdnKeysLoad = singleflight.Group{}
c.invoker = chainMiddlewares(InvokeFunc(c.invokeDirect), c.mw...)
c.tg = tg.NewClient(c.invoker)
}