1319 lines
45 KiB
Go
1319 lines
45 KiB
Go
// Code generated by sqlc. DO NOT EDIT.
|
||
// versions:
|
||
// sqlc v1.31.1
|
||
// source: user_update_event.sql
|
||
|
||
package sqlcgen
|
||
|
||
import (
|
||
"context"
|
||
)
|
||
|
||
const appendUserUpdateEvent = `-- name: AppendUserUpdateEvent :exec
|
||
INSERT INTO user_update_events (
|
||
user_id,
|
||
pts,
|
||
pts_count,
|
||
date,
|
||
event_type,
|
||
event_bool,
|
||
event_peers,
|
||
peer_settings,
|
||
message_ids,
|
||
dialog_filter,
|
||
filter_order,
|
||
folder_peers,
|
||
message_box_id,
|
||
peer_type,
|
||
peer_id,
|
||
filter_id,
|
||
max_id,
|
||
still_unread_count,
|
||
tags_enabled
|
||
) VALUES (
|
||
$1,
|
||
$2,
|
||
$3,
|
||
$4,
|
||
$5,
|
||
$6::boolean,
|
||
$7::jsonb,
|
||
$8::jsonb,
|
||
$9::jsonb,
|
||
$10::jsonb,
|
||
$11::jsonb,
|
||
$12::jsonb,
|
||
$13,
|
||
$14::text,
|
||
$15::bigint,
|
||
$16::int,
|
||
$17::int,
|
||
$18::int,
|
||
$19::boolean
|
||
)
|
||
ON CONFLICT (user_id, pts) DO NOTHING
|
||
`
|
||
|
||
type AppendUserUpdateEventParams struct {
|
||
UserID int64
|
||
Pts int32
|
||
PtsCount int32
|
||
Date int32
|
||
EventType string
|
||
EventBool bool
|
||
EventPeers []byte
|
||
PeerSettings []byte
|
||
MessageIds []byte
|
||
DialogFilter []byte
|
||
FilterOrder []byte
|
||
FolderPeers []byte
|
||
MessageBoxID *int32
|
||
PeerType *string
|
||
PeerID *int64
|
||
FilterID int32
|
||
MaxID int32
|
||
StillUnreadCount int32
|
||
TagsEnabled bool
|
||
}
|
||
|
||
func (q *Queries) AppendUserUpdateEvent(ctx context.Context, arg AppendUserUpdateEventParams) error {
|
||
_, err := q.db.Exec(ctx, appendUserUpdateEvent,
|
||
arg.UserID,
|
||
arg.Pts,
|
||
arg.PtsCount,
|
||
arg.Date,
|
||
arg.EventType,
|
||
arg.EventBool,
|
||
arg.EventPeers,
|
||
arg.PeerSettings,
|
||
arg.MessageIds,
|
||
arg.DialogFilter,
|
||
arg.FilterOrder,
|
||
arg.FolderPeers,
|
||
arg.MessageBoxID,
|
||
arg.PeerType,
|
||
arg.PeerID,
|
||
arg.FilterID,
|
||
arg.MaxID,
|
||
arg.StillUnreadCount,
|
||
arg.TagsEnabled,
|
||
)
|
||
return err
|
||
}
|
||
|
||
const batchListDispatchEvents = `-- name: BatchListDispatchEvents :many
|
||
SELECT
|
||
e.user_id,
|
||
e.pts,
|
||
e.pts_count,
|
||
e.date,
|
||
e.event_type,
|
||
e.event_bool,
|
||
COALESCE(e.event_peers::text, '[]')::text AS event_peers_json,
|
||
COALESCE(e.peer_settings::text, '{}')::text AS peer_settings_json,
|
||
COALESCE(e.message_ids::text, '[]')::text AS message_ids_json,
|
||
COALESCE(e.dialog_filter::text, '{}')::text AS dialog_filter_json,
|
||
COALESCE(e.filter_order::text, '[]')::text AS filter_order_json,
|
||
COALESCE(e.folder_peers::text, '[]')::text AS folder_peers_json,
|
||
COALESCE(e.peer_type, '')::text AS event_peer_type,
|
||
COALESCE(e.peer_id, 0)::bigint AS event_peer_id,
|
||
e.filter_id,
|
||
e.max_id,
|
||
e.still_unread_count,
|
||
e.tags_enabled,
|
||
COALESCE(m.box_id, 0)::int AS message_id,
|
||
COALESCE(m.private_message_id, 0)::bigint AS private_message_id,
|
||
COALESCE(m.owner_user_id, 0)::bigint AS owner_user_id,
|
||
COALESCE(m.peer_type, '')::text AS peer_type,
|
||
COALESCE(m.peer_id, 0)::bigint AS peer_id,
|
||
COALESCE(m.from_user_id, 0)::bigint AS from_user_id,
|
||
COALESCE(m.message_date, 0)::int AS message_date,
|
||
COALESCE(m.edit_date, 0)::int AS edit_date,
|
||
COALESCE(m.outgoing, false)::boolean AS outgoing,
|
||
COALESCE(m.body, '')::text AS body,
|
||
COALESCE(m.entities::text, '[]')::text AS message_entities_json,
|
||
COALESCE(m.silent, false)::boolean AS silent,
|
||
COALESCE(m.noforwards, false)::boolean AS noforwards,
|
||
COALESCE(m.reply_to_msg_id, 0)::int AS reply_to_msg_id,
|
||
COALESCE(m.reply_to_peer_type, '')::text AS reply_to_peer_type,
|
||
COALESCE(m.reply_to_peer_id, 0)::bigint AS reply_to_peer_id,
|
||
COALESCE(m.reply_to_top_id, 0)::int AS reply_to_top_id,
|
||
COALESCE(m.quote_text, '')::text AS quote_text,
|
||
COALESCE(m.quote_entities::text, '[]')::text AS quote_entities_json,
|
||
COALESCE(m.quote_offset, 0)::int AS quote_offset,
|
||
COALESCE(m.fwd_from_peer_type, '')::text AS fwd_from_peer_type,
|
||
COALESCE(m.fwd_from_peer_id, 0)::bigint AS fwd_from_peer_id,
|
||
COALESCE(m.fwd_from_name, '')::text AS fwd_from_name,
|
||
COALESCE(m.fwd_date, 0)::int AS fwd_date,
|
||
COALESCE(m.media::text, '{}')::text AS media_json,
|
||
COALESCE(m.media_unread, false)::boolean AS media_unread,
|
||
COALESCE(m.reaction_unread, false)::boolean AS reaction_unread,
|
||
COALESCE(peer_u.id, 0)::bigint AS peer_user_id,
|
||
COALESCE(peer_u.access_hash, 0)::bigint AS peer_access_hash,
|
||
COALESCE(peer_u.phone, '')::text AS peer_phone,
|
||
COALESCE(peer_u.first_name, '')::text AS peer_first_name,
|
||
COALESCE(peer_u.last_name, '')::text AS peer_last_name,
|
||
COALESCE(peer_u.username, '')::text AS peer_username,
|
||
COALESCE(peer_u.country_code, '')::text AS peer_country_code,
|
||
COALESCE(peer_u.verified, false)::boolean AS peer_verified,
|
||
COALESCE(peer_u.support, false)::boolean AS peer_support,
|
||
COALESCE(from_u.id, 0)::bigint AS from_user_user_id,
|
||
COALESCE(from_u.access_hash, 0)::bigint AS from_user_access_hash,
|
||
COALESCE(from_u.phone, '')::text AS from_user_phone,
|
||
COALESCE(from_u.first_name, '')::text AS from_user_first_name,
|
||
COALESCE(from_u.last_name, '')::text AS from_user_last_name,
|
||
COALESCE(from_u.username, '')::text AS from_user_username,
|
||
COALESCE(from_u.country_code, '')::text AS from_user_country_code,
|
||
COALESCE(from_u.verified, false)::boolean AS from_user_verified,
|
||
COALESCE(from_u.support, false)::boolean AS from_user_support,
|
||
COALESCE(fwd_u.id, 0)::bigint AS fwd_user_id,
|
||
COALESCE(fwd_u.access_hash, 0)::bigint AS fwd_user_access_hash,
|
||
COALESCE(fwd_u.phone, '')::text AS fwd_user_phone,
|
||
COALESCE(fwd_u.first_name, '')::text AS fwd_user_first_name,
|
||
COALESCE(fwd_u.last_name, '')::text AS fwd_user_last_name,
|
||
COALESCE(fwd_u.username, '')::text AS fwd_user_username,
|
||
COALESCE(fwd_u.country_code, '')::text AS fwd_user_country_code,
|
||
COALESCE(fwd_u.verified, false)::boolean AS fwd_user_verified,
|
||
COALESCE(fwd_u.support, false)::boolean AS fwd_user_support,
|
||
COALESCE(reply_u.id, 0)::bigint AS reply_user_id,
|
||
COALESCE(reply_u.access_hash, 0)::bigint AS reply_user_access_hash,
|
||
COALESCE(reply_u.phone, '')::text AS reply_user_phone,
|
||
COALESCE(reply_u.first_name, '')::text AS reply_user_first_name,
|
||
COALESCE(reply_u.last_name, '')::text AS reply_user_last_name,
|
||
COALESCE(reply_u.username, '')::text AS reply_user_username,
|
||
COALESCE(reply_u.country_code, '')::text AS reply_user_country_code,
|
||
COALESCE(reply_u.verified, false)::boolean AS reply_user_verified,
|
||
COALESCE(reply_u.support, false)::boolean AS reply_user_support,
|
||
COALESCE(fwd_ch.id, 0)::bigint AS fwd_channel_id,
|
||
COALESCE(fwd_ch.access_hash, 0)::bigint AS fwd_channel_access_hash,
|
||
COALESCE(fwd_ch.creator_user_id, 0)::bigint AS fwd_channel_creator_user_id,
|
||
COALESCE(fwd_ch.title, '')::text AS fwd_channel_title,
|
||
COALESCE(fwd_ch.about, '')::text AS fwd_channel_about,
|
||
COALESCE(fwd_ch.username, '')::text AS fwd_channel_username,
|
||
COALESCE(fwd_ch.broadcast, false)::boolean AS fwd_channel_broadcast,
|
||
COALESCE(fwd_ch.megagroup, false)::boolean AS fwd_channel_megagroup,
|
||
COALESCE(fwd_ch.forum, false)::boolean AS fwd_channel_forum,
|
||
COALESCE(fwd_ch.noforwards, false)::boolean AS fwd_channel_noforwards,
|
||
COALESCE(fwd_ch.signatures, false)::boolean AS fwd_channel_signatures,
|
||
COALESCE(fwd_ch.pre_history_hidden, false)::boolean AS fwd_channel_pre_history_hidden,
|
||
COALESCE(fwd_ch.slowmode_seconds, 0)::int AS fwd_channel_slowmode_seconds,
|
||
COALESCE(fwd_ch.default_banned_rights::text, '{}')::text AS fwd_channel_default_banned_rights,
|
||
COALESCE(fwd_ch.participants_count, 0)::int AS fwd_channel_participants_count,
|
||
COALESCE(fwd_ch.admins_count, 0)::int AS fwd_channel_admins_count,
|
||
COALESCE(fwd_ch.kicked_count, 0)::int AS fwd_channel_kicked_count,
|
||
COALESCE(fwd_ch.banned_count, 0)::int AS fwd_channel_banned_count,
|
||
COALESCE(fwd_ch.top_message_id, 0)::int AS fwd_channel_top_message_id,
|
||
COALESCE(fwd_ch.pinned_message_id, 0)::int AS fwd_channel_pinned_message_id,
|
||
COALESCE(fwd_ch.pts, 0)::int AS fwd_channel_pts,
|
||
COALESCE(fwd_ch.ttl_period, 0)::int AS fwd_channel_ttl_period,
|
||
COALESCE(fwd_ch.date, 0)::int AS fwd_channel_date,
|
||
COALESCE(fwd_ch.deleted, false)::boolean AS fwd_channel_deleted,
|
||
COALESCE(reply_ch.id, 0)::bigint AS reply_channel_id,
|
||
COALESCE(reply_ch.access_hash, 0)::bigint AS reply_channel_access_hash,
|
||
COALESCE(reply_ch.creator_user_id, 0)::bigint AS reply_channel_creator_user_id,
|
||
COALESCE(reply_ch.title, '')::text AS reply_channel_title,
|
||
COALESCE(reply_ch.about, '')::text AS reply_channel_about,
|
||
COALESCE(reply_ch.username, '')::text AS reply_channel_username,
|
||
COALESCE(reply_ch.broadcast, false)::boolean AS reply_channel_broadcast,
|
||
COALESCE(reply_ch.megagroup, false)::boolean AS reply_channel_megagroup,
|
||
COALESCE(reply_ch.forum, false)::boolean AS reply_channel_forum,
|
||
COALESCE(reply_ch.noforwards, false)::boolean AS reply_channel_noforwards,
|
||
COALESCE(reply_ch.signatures, false)::boolean AS reply_channel_signatures,
|
||
COALESCE(reply_ch.pre_history_hidden, false)::boolean AS reply_channel_pre_history_hidden,
|
||
COALESCE(reply_ch.slowmode_seconds, 0)::int AS reply_channel_slowmode_seconds,
|
||
COALESCE(reply_ch.default_banned_rights::text, '{}')::text AS reply_channel_default_banned_rights,
|
||
COALESCE(reply_ch.participants_count, 0)::int AS reply_channel_participants_count,
|
||
COALESCE(reply_ch.admins_count, 0)::int AS reply_channel_admins_count,
|
||
COALESCE(reply_ch.kicked_count, 0)::int AS reply_channel_kicked_count,
|
||
COALESCE(reply_ch.banned_count, 0)::int AS reply_channel_banned_count,
|
||
COALESCE(reply_ch.top_message_id, 0)::int AS reply_channel_top_message_id,
|
||
COALESCE(reply_ch.pinned_message_id, 0)::int AS reply_channel_pinned_message_id,
|
||
COALESCE(reply_ch.pts, 0)::int AS reply_channel_pts,
|
||
COALESCE(reply_ch.ttl_period, 0)::int AS reply_channel_ttl_period,
|
||
COALESCE(reply_ch.date, 0)::int AS reply_channel_date,
|
||
COALESCE(reply_ch.deleted, false)::boolean AS reply_channel_deleted
|
||
FROM unnest($1::bigint[]) WITH ORDINALITY AS u(user_id, ord)
|
||
JOIN unnest($2::int[]) WITH ORDINALITY AS p(pts, ord) USING (ord)
|
||
JOIN user_update_events e ON e.user_id = u.user_id AND e.pts = p.pts
|
||
LEFT JOIN message_boxes m ON m.owner_user_id = e.user_id AND m.box_id = e.message_box_id
|
||
LEFT JOIN users peer_u ON m.peer_type = 'user' AND peer_u.id = m.peer_id
|
||
LEFT JOIN users from_u ON from_u.id = m.from_user_id
|
||
LEFT JOIN users fwd_u ON m.fwd_from_peer_type = 'user' AND fwd_u.id = m.fwd_from_peer_id
|
||
LEFT JOIN users reply_u ON m.reply_to_peer_type = 'user' AND reply_u.id = m.reply_to_peer_id
|
||
LEFT JOIN channels fwd_ch ON m.fwd_from_peer_type = 'channel' AND fwd_ch.id = m.fwd_from_peer_id
|
||
LEFT JOIN channels reply_ch ON m.reply_to_peer_type = 'channel' AND reply_ch.id = m.reply_to_peer_id
|
||
`
|
||
|
||
type BatchListDispatchEventsParams struct {
|
||
UserIds []int64
|
||
PtsList []int32
|
||
}
|
||
|
||
type BatchListDispatchEventsRow struct {
|
||
UserID int64
|
||
Pts int32
|
||
PtsCount int32
|
||
Date int32
|
||
EventType string
|
||
EventBool bool
|
||
EventPeersJson string
|
||
PeerSettingsJson string
|
||
MessageIdsJson string
|
||
DialogFilterJson string
|
||
FilterOrderJson string
|
||
FolderPeersJson string
|
||
EventPeerType string
|
||
EventPeerID int64
|
||
FilterID int32
|
||
MaxID int32
|
||
StillUnreadCount int32
|
||
TagsEnabled bool
|
||
MessageID int32
|
||
PrivateMessageID int64
|
||
OwnerUserID int64
|
||
PeerType string
|
||
PeerID int64
|
||
FromUserID int64
|
||
MessageDate int32
|
||
EditDate int32
|
||
Outgoing bool
|
||
Body string
|
||
MessageEntitiesJson string
|
||
Silent bool
|
||
Noforwards bool
|
||
ReplyToMsgID int32
|
||
ReplyToPeerType string
|
||
ReplyToPeerID int64
|
||
ReplyToTopID int32
|
||
QuoteText string
|
||
QuoteEntitiesJson string
|
||
QuoteOffset int32
|
||
FwdFromPeerType string
|
||
FwdFromPeerID int64
|
||
FwdFromName string
|
||
FwdDate int32
|
||
MediaJson string
|
||
MediaUnread bool
|
||
ReactionUnread bool
|
||
PeerUserID int64
|
||
PeerAccessHash int64
|
||
PeerPhone string
|
||
PeerFirstName string
|
||
PeerLastName string
|
||
PeerUsername string
|
||
PeerCountryCode string
|
||
PeerVerified bool
|
||
PeerSupport bool
|
||
FromUserUserID int64
|
||
FromUserAccessHash int64
|
||
FromUserPhone string
|
||
FromUserFirstName string
|
||
FromUserLastName string
|
||
FromUserUsername string
|
||
FromUserCountryCode string
|
||
FromUserVerified bool
|
||
FromUserSupport bool
|
||
FwdUserID int64
|
||
FwdUserAccessHash int64
|
||
FwdUserPhone string
|
||
FwdUserFirstName string
|
||
FwdUserLastName string
|
||
FwdUserUsername string
|
||
FwdUserCountryCode string
|
||
FwdUserVerified bool
|
||
FwdUserSupport bool
|
||
ReplyUserID int64
|
||
ReplyUserAccessHash int64
|
||
ReplyUserPhone string
|
||
ReplyUserFirstName string
|
||
ReplyUserLastName string
|
||
ReplyUserUsername string
|
||
ReplyUserCountryCode string
|
||
ReplyUserVerified bool
|
||
ReplyUserSupport bool
|
||
FwdChannelID int64
|
||
FwdChannelAccessHash int64
|
||
FwdChannelCreatorUserID int64
|
||
FwdChannelTitle string
|
||
FwdChannelAbout string
|
||
FwdChannelUsername string
|
||
FwdChannelBroadcast bool
|
||
FwdChannelMegagroup bool
|
||
FwdChannelForum bool
|
||
FwdChannelNoforwards bool
|
||
FwdChannelSignatures bool
|
||
FwdChannelPreHistoryHidden bool
|
||
FwdChannelSlowmodeSeconds int32
|
||
FwdChannelDefaultBannedRights string
|
||
FwdChannelParticipantsCount int32
|
||
FwdChannelAdminsCount int32
|
||
FwdChannelKickedCount int32
|
||
FwdChannelBannedCount int32
|
||
FwdChannelTopMessageID int32
|
||
FwdChannelPinnedMessageID int32
|
||
FwdChannelPts int32
|
||
FwdChannelTtlPeriod int32
|
||
FwdChannelDate int32
|
||
FwdChannelDeleted bool
|
||
ReplyChannelID int64
|
||
ReplyChannelAccessHash int64
|
||
ReplyChannelCreatorUserID int64
|
||
ReplyChannelTitle string
|
||
ReplyChannelAbout string
|
||
ReplyChannelUsername string
|
||
ReplyChannelBroadcast bool
|
||
ReplyChannelMegagroup bool
|
||
ReplyChannelForum bool
|
||
ReplyChannelNoforwards bool
|
||
ReplyChannelSignatures bool
|
||
ReplyChannelPreHistoryHidden bool
|
||
ReplyChannelSlowmodeSeconds int32
|
||
ReplyChannelDefaultBannedRights string
|
||
ReplyChannelParticipantsCount int32
|
||
ReplyChannelAdminsCount int32
|
||
ReplyChannelKickedCount int32
|
||
ReplyChannelBannedCount int32
|
||
ReplyChannelTopMessageID int32
|
||
ReplyChannelPinnedMessageID int32
|
||
ReplyChannelPts int32
|
||
ReplyChannelTtlPeriod int32
|
||
ReplyChannelDate int32
|
||
ReplyChannelDeleted bool
|
||
}
|
||
|
||
// 按 (user_id, pts) 精确批量取账号事件,供 outbox worker 一次性加载一批 claim 的事件详情,
|
||
// 取代逐条 ListUserUpdateEventsAfter。列与 ListUserUpdateEventsAfter 完全一致以复用转换逻辑。
|
||
func (q *Queries) BatchListDispatchEvents(ctx context.Context, arg BatchListDispatchEventsParams) ([]BatchListDispatchEventsRow, error) {
|
||
rows, err := q.db.Query(ctx, batchListDispatchEvents, arg.UserIds, arg.PtsList)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
var items []BatchListDispatchEventsRow
|
||
for rows.Next() {
|
||
var i BatchListDispatchEventsRow
|
||
if err := rows.Scan(
|
||
&i.UserID,
|
||
&i.Pts,
|
||
&i.PtsCount,
|
||
&i.Date,
|
||
&i.EventType,
|
||
&i.EventBool,
|
||
&i.EventPeersJson,
|
||
&i.PeerSettingsJson,
|
||
&i.MessageIdsJson,
|
||
&i.DialogFilterJson,
|
||
&i.FilterOrderJson,
|
||
&i.FolderPeersJson,
|
||
&i.EventPeerType,
|
||
&i.EventPeerID,
|
||
&i.FilterID,
|
||
&i.MaxID,
|
||
&i.StillUnreadCount,
|
||
&i.TagsEnabled,
|
||
&i.MessageID,
|
||
&i.PrivateMessageID,
|
||
&i.OwnerUserID,
|
||
&i.PeerType,
|
||
&i.PeerID,
|
||
&i.FromUserID,
|
||
&i.MessageDate,
|
||
&i.EditDate,
|
||
&i.Outgoing,
|
||
&i.Body,
|
||
&i.MessageEntitiesJson,
|
||
&i.Silent,
|
||
&i.Noforwards,
|
||
&i.ReplyToMsgID,
|
||
&i.ReplyToPeerType,
|
||
&i.ReplyToPeerID,
|
||
&i.ReplyToTopID,
|
||
&i.QuoteText,
|
||
&i.QuoteEntitiesJson,
|
||
&i.QuoteOffset,
|
||
&i.FwdFromPeerType,
|
||
&i.FwdFromPeerID,
|
||
&i.FwdFromName,
|
||
&i.FwdDate,
|
||
&i.MediaJson,
|
||
&i.MediaUnread,
|
||
&i.ReactionUnread,
|
||
&i.PeerUserID,
|
||
&i.PeerAccessHash,
|
||
&i.PeerPhone,
|
||
&i.PeerFirstName,
|
||
&i.PeerLastName,
|
||
&i.PeerUsername,
|
||
&i.PeerCountryCode,
|
||
&i.PeerVerified,
|
||
&i.PeerSupport,
|
||
&i.FromUserUserID,
|
||
&i.FromUserAccessHash,
|
||
&i.FromUserPhone,
|
||
&i.FromUserFirstName,
|
||
&i.FromUserLastName,
|
||
&i.FromUserUsername,
|
||
&i.FromUserCountryCode,
|
||
&i.FromUserVerified,
|
||
&i.FromUserSupport,
|
||
&i.FwdUserID,
|
||
&i.FwdUserAccessHash,
|
||
&i.FwdUserPhone,
|
||
&i.FwdUserFirstName,
|
||
&i.FwdUserLastName,
|
||
&i.FwdUserUsername,
|
||
&i.FwdUserCountryCode,
|
||
&i.FwdUserVerified,
|
||
&i.FwdUserSupport,
|
||
&i.ReplyUserID,
|
||
&i.ReplyUserAccessHash,
|
||
&i.ReplyUserPhone,
|
||
&i.ReplyUserFirstName,
|
||
&i.ReplyUserLastName,
|
||
&i.ReplyUserUsername,
|
||
&i.ReplyUserCountryCode,
|
||
&i.ReplyUserVerified,
|
||
&i.ReplyUserSupport,
|
||
&i.FwdChannelID,
|
||
&i.FwdChannelAccessHash,
|
||
&i.FwdChannelCreatorUserID,
|
||
&i.FwdChannelTitle,
|
||
&i.FwdChannelAbout,
|
||
&i.FwdChannelUsername,
|
||
&i.FwdChannelBroadcast,
|
||
&i.FwdChannelMegagroup,
|
||
&i.FwdChannelForum,
|
||
&i.FwdChannelNoforwards,
|
||
&i.FwdChannelSignatures,
|
||
&i.FwdChannelPreHistoryHidden,
|
||
&i.FwdChannelSlowmodeSeconds,
|
||
&i.FwdChannelDefaultBannedRights,
|
||
&i.FwdChannelParticipantsCount,
|
||
&i.FwdChannelAdminsCount,
|
||
&i.FwdChannelKickedCount,
|
||
&i.FwdChannelBannedCount,
|
||
&i.FwdChannelTopMessageID,
|
||
&i.FwdChannelPinnedMessageID,
|
||
&i.FwdChannelPts,
|
||
&i.FwdChannelTtlPeriod,
|
||
&i.FwdChannelDate,
|
||
&i.FwdChannelDeleted,
|
||
&i.ReplyChannelID,
|
||
&i.ReplyChannelAccessHash,
|
||
&i.ReplyChannelCreatorUserID,
|
||
&i.ReplyChannelTitle,
|
||
&i.ReplyChannelAbout,
|
||
&i.ReplyChannelUsername,
|
||
&i.ReplyChannelBroadcast,
|
||
&i.ReplyChannelMegagroup,
|
||
&i.ReplyChannelForum,
|
||
&i.ReplyChannelNoforwards,
|
||
&i.ReplyChannelSignatures,
|
||
&i.ReplyChannelPreHistoryHidden,
|
||
&i.ReplyChannelSlowmodeSeconds,
|
||
&i.ReplyChannelDefaultBannedRights,
|
||
&i.ReplyChannelParticipantsCount,
|
||
&i.ReplyChannelAdminsCount,
|
||
&i.ReplyChannelKickedCount,
|
||
&i.ReplyChannelBannedCount,
|
||
&i.ReplyChannelTopMessageID,
|
||
&i.ReplyChannelPinnedMessageID,
|
||
&i.ReplyChannelPts,
|
||
&i.ReplyChannelTtlPeriod,
|
||
&i.ReplyChannelDate,
|
||
&i.ReplyChannelDeleted,
|
||
); err != nil {
|
||
return nil, err
|
||
}
|
||
items = append(items, i)
|
||
}
|
||
if err := rows.Err(); err != nil {
|
||
return nil, err
|
||
}
|
||
return items, nil
|
||
}
|
||
|
||
const claimDispatchOutbox = `-- name: ClaimDispatchOutbox :many
|
||
WITH picked AS (
|
||
SELECT target_user_id, id
|
||
FROM dispatch_outbox
|
||
WHERE (
|
||
status = 'pending'
|
||
AND next_attempt_at <= now()
|
||
)
|
||
OR (
|
||
status = 'dispatching'
|
||
AND updated_at < now() - make_interval(secs => $1::int)
|
||
)
|
||
ORDER BY next_attempt_at ASC, target_user_id ASC, id ASC
|
||
LIMIT $2
|
||
FOR UPDATE SKIP LOCKED
|
||
)
|
||
UPDATE dispatch_outbox d
|
||
SET
|
||
status = 'dispatching',
|
||
attempts = d.attempts + 1,
|
||
updated_at = now()
|
||
FROM picked p
|
||
WHERE d.target_user_id = p.target_user_id
|
||
AND d.id = p.id
|
||
RETURNING
|
||
d.id,
|
||
d.target_user_id,
|
||
d.pts,
|
||
d.event_type,
|
||
d.exclude_auth_key_id,
|
||
d.exclude_session_id,
|
||
d.attempts
|
||
`
|
||
|
||
type ClaimDispatchOutboxParams struct {
|
||
LeaseSeconds int32
|
||
LimitCount int32
|
||
}
|
||
|
||
type ClaimDispatchOutboxRow struct {
|
||
ID int64
|
||
TargetUserID int64
|
||
Pts int32
|
||
EventType string
|
||
ExcludeAuthKeyID int64
|
||
ExcludeSessionID int64
|
||
Attempts int32
|
||
}
|
||
|
||
func (q *Queries) ClaimDispatchOutbox(ctx context.Context, arg ClaimDispatchOutboxParams) ([]ClaimDispatchOutboxRow, error) {
|
||
rows, err := q.db.Query(ctx, claimDispatchOutbox, arg.LeaseSeconds, arg.LimitCount)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
var items []ClaimDispatchOutboxRow
|
||
for rows.Next() {
|
||
var i ClaimDispatchOutboxRow
|
||
if err := rows.Scan(
|
||
&i.ID,
|
||
&i.TargetUserID,
|
||
&i.Pts,
|
||
&i.EventType,
|
||
&i.ExcludeAuthKeyID,
|
||
&i.ExcludeSessionID,
|
||
&i.Attempts,
|
||
); err != nil {
|
||
return nil, err
|
||
}
|
||
items = append(items, i)
|
||
}
|
||
if err := rows.Err(); err != nil {
|
||
return nil, err
|
||
}
|
||
return items, nil
|
||
}
|
||
|
||
const deleteFailedDispatchOutbox = `-- name: DeleteFailedDispatchOutbox :one
|
||
WITH doomed AS (
|
||
SELECT target_user_id, id
|
||
FROM dispatch_outbox
|
||
WHERE status = 'failed'
|
||
AND updated_at < now() - make_interval(secs => $1::int)
|
||
ORDER BY updated_at ASC, target_user_id ASC, id ASC
|
||
LIMIT $2
|
||
),
|
||
deleted AS (
|
||
DELETE FROM dispatch_outbox d
|
||
USING doomed x
|
||
WHERE d.target_user_id = x.target_user_id
|
||
AND d.id = x.id
|
||
RETURNING d.id
|
||
)
|
||
SELECT count(*)::int AS deleted_count
|
||
FROM deleted
|
||
`
|
||
|
||
type DeleteFailedDispatchOutboxParams struct {
|
||
OlderThanSeconds int32
|
||
LimitCount int32
|
||
}
|
||
|
||
func (q *Queries) DeleteFailedDispatchOutbox(ctx context.Context, arg DeleteFailedDispatchOutboxParams) (int32, error) {
|
||
row := q.db.QueryRow(ctx, deleteFailedDispatchOutbox, arg.OlderThanSeconds, arg.LimitCount)
|
||
var deleted_count int32
|
||
err := row.Scan(&deleted_count)
|
||
return deleted_count, err
|
||
}
|
||
|
||
const enqueueDispatch = `-- name: EnqueueDispatch :exec
|
||
INSERT INTO dispatch_outbox (
|
||
target_user_id,
|
||
pts,
|
||
event_type,
|
||
exclude_auth_key_id,
|
||
exclude_session_id
|
||
) VALUES (
|
||
$1, $2, $3, $4, $5
|
||
)
|
||
ON CONFLICT DO NOTHING
|
||
`
|
||
|
||
type EnqueueDispatchParams struct {
|
||
TargetUserID int64
|
||
Pts int32
|
||
EventType string
|
||
ExcludeAuthKeyID int64
|
||
ExcludeSessionID int64
|
||
}
|
||
|
||
func (q *Queries) EnqueueDispatch(ctx context.Context, arg EnqueueDispatchParams) error {
|
||
_, err := q.db.Exec(ctx, enqueueDispatch,
|
||
arg.TargetUserID,
|
||
arg.Pts,
|
||
arg.EventType,
|
||
arg.ExcludeAuthKeyID,
|
||
arg.ExcludeSessionID,
|
||
)
|
||
return err
|
||
}
|
||
|
||
const ensureUserUpdateWatermark = `-- name: EnsureUserUpdateWatermark :exec
|
||
INSERT INTO user_update_watermarks (user_id, contiguous_pts)
|
||
VALUES ($1, 0)
|
||
ON CONFLICT (user_id) DO NOTHING
|
||
`
|
||
|
||
func (q *Queries) EnsureUserUpdateWatermark(ctx context.Context, userID int64) error {
|
||
_, err := q.db.Exec(ctx, ensureUserUpdateWatermark, userID)
|
||
return err
|
||
}
|
||
|
||
const getUserUpdateWatermark = `-- name: GetUserUpdateWatermark :one
|
||
SELECT contiguous_pts
|
||
FROM user_update_watermarks
|
||
WHERE user_id = $1
|
||
`
|
||
|
||
func (q *Queries) GetUserUpdateWatermark(ctx context.Context, userID int64) (int32, error) {
|
||
row := q.db.QueryRow(ctx, getUserUpdateWatermark, userID)
|
||
var contiguous_pts int32
|
||
err := row.Scan(&contiguous_pts)
|
||
return contiguous_pts, err
|
||
}
|
||
|
||
const listUserUpdateEventsAfter = `-- name: ListUserUpdateEventsAfter :many
|
||
SELECT
|
||
e.user_id,
|
||
e.pts,
|
||
e.pts_count,
|
||
e.date,
|
||
e.event_type,
|
||
e.event_bool,
|
||
COALESCE(e.event_peers::text, '[]')::text AS event_peers_json,
|
||
COALESCE(e.peer_settings::text, '{}')::text AS peer_settings_json,
|
||
COALESCE(e.message_ids::text, '[]')::text AS message_ids_json,
|
||
COALESCE(e.dialog_filter::text, '{}')::text AS dialog_filter_json,
|
||
COALESCE(e.filter_order::text, '[]')::text AS filter_order_json,
|
||
COALESCE(e.folder_peers::text, '[]')::text AS folder_peers_json,
|
||
COALESCE(e.peer_type, '')::text AS event_peer_type,
|
||
COALESCE(e.peer_id, 0)::bigint AS event_peer_id,
|
||
e.filter_id,
|
||
e.max_id,
|
||
e.still_unread_count,
|
||
e.tags_enabled,
|
||
COALESCE(m.box_id, 0)::int AS message_id,
|
||
COALESCE(m.private_message_id, 0)::bigint AS private_message_id,
|
||
COALESCE(m.owner_user_id, 0)::bigint AS owner_user_id,
|
||
COALESCE(m.peer_type, '')::text AS peer_type,
|
||
COALESCE(m.peer_id, 0)::bigint AS peer_id,
|
||
COALESCE(m.from_user_id, 0)::bigint AS from_user_id,
|
||
COALESCE(m.message_date, 0)::int AS message_date,
|
||
COALESCE(m.edit_date, 0)::int AS edit_date,
|
||
COALESCE(m.outgoing, false)::boolean AS outgoing,
|
||
COALESCE(m.body, '')::text AS body,
|
||
COALESCE(m.entities::text, '[]')::text AS message_entities_json,
|
||
COALESCE(m.silent, false)::boolean AS silent,
|
||
COALESCE(m.noforwards, false)::boolean AS noforwards,
|
||
COALESCE(m.reply_to_msg_id, 0)::int AS reply_to_msg_id,
|
||
COALESCE(m.reply_to_peer_type, '')::text AS reply_to_peer_type,
|
||
COALESCE(m.reply_to_peer_id, 0)::bigint AS reply_to_peer_id,
|
||
COALESCE(m.reply_to_top_id, 0)::int AS reply_to_top_id,
|
||
COALESCE(m.quote_text, '')::text AS quote_text,
|
||
COALESCE(m.quote_entities::text, '[]')::text AS quote_entities_json,
|
||
COALESCE(m.quote_offset, 0)::int AS quote_offset,
|
||
COALESCE(m.fwd_from_peer_type, '')::text AS fwd_from_peer_type,
|
||
COALESCE(m.fwd_from_peer_id, 0)::bigint AS fwd_from_peer_id,
|
||
COALESCE(m.fwd_from_name, '')::text AS fwd_from_name,
|
||
COALESCE(m.fwd_date, 0)::int AS fwd_date,
|
||
COALESCE(m.media::text, '{}')::text AS media_json,
|
||
COALESCE(m.media_unread, false)::boolean AS media_unread,
|
||
COALESCE(m.reaction_unread, false)::boolean AS reaction_unread,
|
||
COALESCE(peer_u.id, 0)::bigint AS peer_user_id,
|
||
COALESCE(peer_u.access_hash, 0)::bigint AS peer_access_hash,
|
||
COALESCE(peer_u.phone, '')::text AS peer_phone,
|
||
COALESCE(peer_u.first_name, '')::text AS peer_first_name,
|
||
COALESCE(peer_u.last_name, '')::text AS peer_last_name,
|
||
COALESCE(peer_u.username, '')::text AS peer_username,
|
||
COALESCE(peer_u.country_code, '')::text AS peer_country_code,
|
||
COALESCE(peer_u.verified, false)::boolean AS peer_verified,
|
||
COALESCE(peer_u.support, false)::boolean AS peer_support,
|
||
COALESCE(from_u.id, 0)::bigint AS from_user_user_id,
|
||
COALESCE(from_u.access_hash, 0)::bigint AS from_user_access_hash,
|
||
COALESCE(from_u.phone, '')::text AS from_user_phone,
|
||
COALESCE(from_u.first_name, '')::text AS from_user_first_name,
|
||
COALESCE(from_u.last_name, '')::text AS from_user_last_name,
|
||
COALESCE(from_u.username, '')::text AS from_user_username,
|
||
COALESCE(from_u.country_code, '')::text AS from_user_country_code,
|
||
COALESCE(from_u.verified, false)::boolean AS from_user_verified,
|
||
COALESCE(from_u.support, false)::boolean AS from_user_support,
|
||
COALESCE(fwd_u.id, 0)::bigint AS fwd_user_id,
|
||
COALESCE(fwd_u.access_hash, 0)::bigint AS fwd_user_access_hash,
|
||
COALESCE(fwd_u.phone, '')::text AS fwd_user_phone,
|
||
COALESCE(fwd_u.first_name, '')::text AS fwd_user_first_name,
|
||
COALESCE(fwd_u.last_name, '')::text AS fwd_user_last_name,
|
||
COALESCE(fwd_u.username, '')::text AS fwd_user_username,
|
||
COALESCE(fwd_u.country_code, '')::text AS fwd_user_country_code,
|
||
COALESCE(fwd_u.verified, false)::boolean AS fwd_user_verified,
|
||
COALESCE(fwd_u.support, false)::boolean AS fwd_user_support,
|
||
COALESCE(reply_u.id, 0)::bigint AS reply_user_id,
|
||
COALESCE(reply_u.access_hash, 0)::bigint AS reply_user_access_hash,
|
||
COALESCE(reply_u.phone, '')::text AS reply_user_phone,
|
||
COALESCE(reply_u.first_name, '')::text AS reply_user_first_name,
|
||
COALESCE(reply_u.last_name, '')::text AS reply_user_last_name,
|
||
COALESCE(reply_u.username, '')::text AS reply_user_username,
|
||
COALESCE(reply_u.country_code, '')::text AS reply_user_country_code,
|
||
COALESCE(reply_u.verified, false)::boolean AS reply_user_verified,
|
||
COALESCE(reply_u.support, false)::boolean AS reply_user_support,
|
||
COALESCE(fwd_ch.id, 0)::bigint AS fwd_channel_id,
|
||
COALESCE(fwd_ch.access_hash, 0)::bigint AS fwd_channel_access_hash,
|
||
COALESCE(fwd_ch.creator_user_id, 0)::bigint AS fwd_channel_creator_user_id,
|
||
COALESCE(fwd_ch.title, '')::text AS fwd_channel_title,
|
||
COALESCE(fwd_ch.about, '')::text AS fwd_channel_about,
|
||
COALESCE(fwd_ch.username, '')::text AS fwd_channel_username,
|
||
COALESCE(fwd_ch.broadcast, false)::boolean AS fwd_channel_broadcast,
|
||
COALESCE(fwd_ch.megagroup, false)::boolean AS fwd_channel_megagroup,
|
||
COALESCE(fwd_ch.forum, false)::boolean AS fwd_channel_forum,
|
||
COALESCE(fwd_ch.noforwards, false)::boolean AS fwd_channel_noforwards,
|
||
COALESCE(fwd_ch.signatures, false)::boolean AS fwd_channel_signatures,
|
||
COALESCE(fwd_ch.pre_history_hidden, false)::boolean AS fwd_channel_pre_history_hidden,
|
||
COALESCE(fwd_ch.slowmode_seconds, 0)::int AS fwd_channel_slowmode_seconds,
|
||
COALESCE(fwd_ch.default_banned_rights::text, '{}')::text AS fwd_channel_default_banned_rights,
|
||
COALESCE(fwd_ch.participants_count, 0)::int AS fwd_channel_participants_count,
|
||
COALESCE(fwd_ch.admins_count, 0)::int AS fwd_channel_admins_count,
|
||
COALESCE(fwd_ch.kicked_count, 0)::int AS fwd_channel_kicked_count,
|
||
COALESCE(fwd_ch.banned_count, 0)::int AS fwd_channel_banned_count,
|
||
COALESCE(fwd_ch.top_message_id, 0)::int AS fwd_channel_top_message_id,
|
||
COALESCE(fwd_ch.pinned_message_id, 0)::int AS fwd_channel_pinned_message_id,
|
||
COALESCE(fwd_ch.pts, 0)::int AS fwd_channel_pts,
|
||
COALESCE(fwd_ch.ttl_period, 0)::int AS fwd_channel_ttl_period,
|
||
COALESCE(fwd_ch.date, 0)::int AS fwd_channel_date,
|
||
COALESCE(fwd_ch.deleted, false)::boolean AS fwd_channel_deleted,
|
||
COALESCE(reply_ch.id, 0)::bigint AS reply_channel_id,
|
||
COALESCE(reply_ch.access_hash, 0)::bigint AS reply_channel_access_hash,
|
||
COALESCE(reply_ch.creator_user_id, 0)::bigint AS reply_channel_creator_user_id,
|
||
COALESCE(reply_ch.title, '')::text AS reply_channel_title,
|
||
COALESCE(reply_ch.about, '')::text AS reply_channel_about,
|
||
COALESCE(reply_ch.username, '')::text AS reply_channel_username,
|
||
COALESCE(reply_ch.broadcast, false)::boolean AS reply_channel_broadcast,
|
||
COALESCE(reply_ch.megagroup, false)::boolean AS reply_channel_megagroup,
|
||
COALESCE(reply_ch.forum, false)::boolean AS reply_channel_forum,
|
||
COALESCE(reply_ch.noforwards, false)::boolean AS reply_channel_noforwards,
|
||
COALESCE(reply_ch.signatures, false)::boolean AS reply_channel_signatures,
|
||
COALESCE(reply_ch.pre_history_hidden, false)::boolean AS reply_channel_pre_history_hidden,
|
||
COALESCE(reply_ch.slowmode_seconds, 0)::int AS reply_channel_slowmode_seconds,
|
||
COALESCE(reply_ch.default_banned_rights::text, '{}')::text AS reply_channel_default_banned_rights,
|
||
COALESCE(reply_ch.participants_count, 0)::int AS reply_channel_participants_count,
|
||
COALESCE(reply_ch.admins_count, 0)::int AS reply_channel_admins_count,
|
||
COALESCE(reply_ch.kicked_count, 0)::int AS reply_channel_kicked_count,
|
||
COALESCE(reply_ch.banned_count, 0)::int AS reply_channel_banned_count,
|
||
COALESCE(reply_ch.top_message_id, 0)::int AS reply_channel_top_message_id,
|
||
COALESCE(reply_ch.pinned_message_id, 0)::int AS reply_channel_pinned_message_id,
|
||
COALESCE(reply_ch.pts, 0)::int AS reply_channel_pts,
|
||
COALESCE(reply_ch.ttl_period, 0)::int AS reply_channel_ttl_period,
|
||
COALESCE(reply_ch.date, 0)::int AS reply_channel_date,
|
||
COALESCE(reply_ch.deleted, false)::boolean AS reply_channel_deleted
|
||
FROM user_update_events e
|
||
LEFT JOIN message_boxes m ON m.owner_user_id = e.user_id AND m.box_id = e.message_box_id
|
||
LEFT JOIN users peer_u ON m.peer_type = 'user' AND peer_u.id = m.peer_id
|
||
LEFT JOIN users from_u ON from_u.id = m.from_user_id
|
||
LEFT JOIN users fwd_u ON m.fwd_from_peer_type = 'user' AND fwd_u.id = m.fwd_from_peer_id
|
||
LEFT JOIN users reply_u ON m.reply_to_peer_type = 'user' AND reply_u.id = m.reply_to_peer_id
|
||
LEFT JOIN channels fwd_ch ON m.fwd_from_peer_type = 'channel' AND fwd_ch.id = m.fwd_from_peer_id
|
||
LEFT JOIN channels reply_ch ON m.reply_to_peer_type = 'channel' AND reply_ch.id = m.reply_to_peer_id
|
||
WHERE e.user_id = $1
|
||
AND e.pts > $2
|
||
ORDER BY e.pts ASC
|
||
LIMIT $3
|
||
`
|
||
|
||
type ListUserUpdateEventsAfterParams struct {
|
||
UserID int64
|
||
Pts int32
|
||
LimitCount int32
|
||
}
|
||
|
||
type ListUserUpdateEventsAfterRow struct {
|
||
UserID int64
|
||
Pts int32
|
||
PtsCount int32
|
||
Date int32
|
||
EventType string
|
||
EventBool bool
|
||
EventPeersJson string
|
||
PeerSettingsJson string
|
||
MessageIdsJson string
|
||
DialogFilterJson string
|
||
FilterOrderJson string
|
||
FolderPeersJson string
|
||
EventPeerType string
|
||
EventPeerID int64
|
||
FilterID int32
|
||
MaxID int32
|
||
StillUnreadCount int32
|
||
TagsEnabled bool
|
||
MessageID int32
|
||
PrivateMessageID int64
|
||
OwnerUserID int64
|
||
PeerType string
|
||
PeerID int64
|
||
FromUserID int64
|
||
MessageDate int32
|
||
EditDate int32
|
||
Outgoing bool
|
||
Body string
|
||
MessageEntitiesJson string
|
||
Silent bool
|
||
Noforwards bool
|
||
ReplyToMsgID int32
|
||
ReplyToPeerType string
|
||
ReplyToPeerID int64
|
||
ReplyToTopID int32
|
||
QuoteText string
|
||
QuoteEntitiesJson string
|
||
QuoteOffset int32
|
||
FwdFromPeerType string
|
||
FwdFromPeerID int64
|
||
FwdFromName string
|
||
FwdDate int32
|
||
MediaJson string
|
||
MediaUnread bool
|
||
ReactionUnread bool
|
||
PeerUserID int64
|
||
PeerAccessHash int64
|
||
PeerPhone string
|
||
PeerFirstName string
|
||
PeerLastName string
|
||
PeerUsername string
|
||
PeerCountryCode string
|
||
PeerVerified bool
|
||
PeerSupport bool
|
||
FromUserUserID int64
|
||
FromUserAccessHash int64
|
||
FromUserPhone string
|
||
FromUserFirstName string
|
||
FromUserLastName string
|
||
FromUserUsername string
|
||
FromUserCountryCode string
|
||
FromUserVerified bool
|
||
FromUserSupport bool
|
||
FwdUserID int64
|
||
FwdUserAccessHash int64
|
||
FwdUserPhone string
|
||
FwdUserFirstName string
|
||
FwdUserLastName string
|
||
FwdUserUsername string
|
||
FwdUserCountryCode string
|
||
FwdUserVerified bool
|
||
FwdUserSupport bool
|
||
ReplyUserID int64
|
||
ReplyUserAccessHash int64
|
||
ReplyUserPhone string
|
||
ReplyUserFirstName string
|
||
ReplyUserLastName string
|
||
ReplyUserUsername string
|
||
ReplyUserCountryCode string
|
||
ReplyUserVerified bool
|
||
ReplyUserSupport bool
|
||
FwdChannelID int64
|
||
FwdChannelAccessHash int64
|
||
FwdChannelCreatorUserID int64
|
||
FwdChannelTitle string
|
||
FwdChannelAbout string
|
||
FwdChannelUsername string
|
||
FwdChannelBroadcast bool
|
||
FwdChannelMegagroup bool
|
||
FwdChannelForum bool
|
||
FwdChannelNoforwards bool
|
||
FwdChannelSignatures bool
|
||
FwdChannelPreHistoryHidden bool
|
||
FwdChannelSlowmodeSeconds int32
|
||
FwdChannelDefaultBannedRights string
|
||
FwdChannelParticipantsCount int32
|
||
FwdChannelAdminsCount int32
|
||
FwdChannelKickedCount int32
|
||
FwdChannelBannedCount int32
|
||
FwdChannelTopMessageID int32
|
||
FwdChannelPinnedMessageID int32
|
||
FwdChannelPts int32
|
||
FwdChannelTtlPeriod int32
|
||
FwdChannelDate int32
|
||
FwdChannelDeleted bool
|
||
ReplyChannelID int64
|
||
ReplyChannelAccessHash int64
|
||
ReplyChannelCreatorUserID int64
|
||
ReplyChannelTitle string
|
||
ReplyChannelAbout string
|
||
ReplyChannelUsername string
|
||
ReplyChannelBroadcast bool
|
||
ReplyChannelMegagroup bool
|
||
ReplyChannelForum bool
|
||
ReplyChannelNoforwards bool
|
||
ReplyChannelSignatures bool
|
||
ReplyChannelPreHistoryHidden bool
|
||
ReplyChannelSlowmodeSeconds int32
|
||
ReplyChannelDefaultBannedRights string
|
||
ReplyChannelParticipantsCount int32
|
||
ReplyChannelAdminsCount int32
|
||
ReplyChannelKickedCount int32
|
||
ReplyChannelBannedCount int32
|
||
ReplyChannelTopMessageID int32
|
||
ReplyChannelPinnedMessageID int32
|
||
ReplyChannelPts int32
|
||
ReplyChannelTtlPeriod int32
|
||
ReplyChannelDate int32
|
||
ReplyChannelDeleted bool
|
||
}
|
||
|
||
func (q *Queries) ListUserUpdateEventsAfter(ctx context.Context, arg ListUserUpdateEventsAfterParams) ([]ListUserUpdateEventsAfterRow, error) {
|
||
rows, err := q.db.Query(ctx, listUserUpdateEventsAfter, arg.UserID, arg.Pts, arg.LimitCount)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
var items []ListUserUpdateEventsAfterRow
|
||
for rows.Next() {
|
||
var i ListUserUpdateEventsAfterRow
|
||
if err := rows.Scan(
|
||
&i.UserID,
|
||
&i.Pts,
|
||
&i.PtsCount,
|
||
&i.Date,
|
||
&i.EventType,
|
||
&i.EventBool,
|
||
&i.EventPeersJson,
|
||
&i.PeerSettingsJson,
|
||
&i.MessageIdsJson,
|
||
&i.DialogFilterJson,
|
||
&i.FilterOrderJson,
|
||
&i.FolderPeersJson,
|
||
&i.EventPeerType,
|
||
&i.EventPeerID,
|
||
&i.FilterID,
|
||
&i.MaxID,
|
||
&i.StillUnreadCount,
|
||
&i.TagsEnabled,
|
||
&i.MessageID,
|
||
&i.PrivateMessageID,
|
||
&i.OwnerUserID,
|
||
&i.PeerType,
|
||
&i.PeerID,
|
||
&i.FromUserID,
|
||
&i.MessageDate,
|
||
&i.EditDate,
|
||
&i.Outgoing,
|
||
&i.Body,
|
||
&i.MessageEntitiesJson,
|
||
&i.Silent,
|
||
&i.Noforwards,
|
||
&i.ReplyToMsgID,
|
||
&i.ReplyToPeerType,
|
||
&i.ReplyToPeerID,
|
||
&i.ReplyToTopID,
|
||
&i.QuoteText,
|
||
&i.QuoteEntitiesJson,
|
||
&i.QuoteOffset,
|
||
&i.FwdFromPeerType,
|
||
&i.FwdFromPeerID,
|
||
&i.FwdFromName,
|
||
&i.FwdDate,
|
||
&i.MediaJson,
|
||
&i.MediaUnread,
|
||
&i.ReactionUnread,
|
||
&i.PeerUserID,
|
||
&i.PeerAccessHash,
|
||
&i.PeerPhone,
|
||
&i.PeerFirstName,
|
||
&i.PeerLastName,
|
||
&i.PeerUsername,
|
||
&i.PeerCountryCode,
|
||
&i.PeerVerified,
|
||
&i.PeerSupport,
|
||
&i.FromUserUserID,
|
||
&i.FromUserAccessHash,
|
||
&i.FromUserPhone,
|
||
&i.FromUserFirstName,
|
||
&i.FromUserLastName,
|
||
&i.FromUserUsername,
|
||
&i.FromUserCountryCode,
|
||
&i.FromUserVerified,
|
||
&i.FromUserSupport,
|
||
&i.FwdUserID,
|
||
&i.FwdUserAccessHash,
|
||
&i.FwdUserPhone,
|
||
&i.FwdUserFirstName,
|
||
&i.FwdUserLastName,
|
||
&i.FwdUserUsername,
|
||
&i.FwdUserCountryCode,
|
||
&i.FwdUserVerified,
|
||
&i.FwdUserSupport,
|
||
&i.ReplyUserID,
|
||
&i.ReplyUserAccessHash,
|
||
&i.ReplyUserPhone,
|
||
&i.ReplyUserFirstName,
|
||
&i.ReplyUserLastName,
|
||
&i.ReplyUserUsername,
|
||
&i.ReplyUserCountryCode,
|
||
&i.ReplyUserVerified,
|
||
&i.ReplyUserSupport,
|
||
&i.FwdChannelID,
|
||
&i.FwdChannelAccessHash,
|
||
&i.FwdChannelCreatorUserID,
|
||
&i.FwdChannelTitle,
|
||
&i.FwdChannelAbout,
|
||
&i.FwdChannelUsername,
|
||
&i.FwdChannelBroadcast,
|
||
&i.FwdChannelMegagroup,
|
||
&i.FwdChannelForum,
|
||
&i.FwdChannelNoforwards,
|
||
&i.FwdChannelSignatures,
|
||
&i.FwdChannelPreHistoryHidden,
|
||
&i.FwdChannelSlowmodeSeconds,
|
||
&i.FwdChannelDefaultBannedRights,
|
||
&i.FwdChannelParticipantsCount,
|
||
&i.FwdChannelAdminsCount,
|
||
&i.FwdChannelKickedCount,
|
||
&i.FwdChannelBannedCount,
|
||
&i.FwdChannelTopMessageID,
|
||
&i.FwdChannelPinnedMessageID,
|
||
&i.FwdChannelPts,
|
||
&i.FwdChannelTtlPeriod,
|
||
&i.FwdChannelDate,
|
||
&i.FwdChannelDeleted,
|
||
&i.ReplyChannelID,
|
||
&i.ReplyChannelAccessHash,
|
||
&i.ReplyChannelCreatorUserID,
|
||
&i.ReplyChannelTitle,
|
||
&i.ReplyChannelAbout,
|
||
&i.ReplyChannelUsername,
|
||
&i.ReplyChannelBroadcast,
|
||
&i.ReplyChannelMegagroup,
|
||
&i.ReplyChannelForum,
|
||
&i.ReplyChannelNoforwards,
|
||
&i.ReplyChannelSignatures,
|
||
&i.ReplyChannelPreHistoryHidden,
|
||
&i.ReplyChannelSlowmodeSeconds,
|
||
&i.ReplyChannelDefaultBannedRights,
|
||
&i.ReplyChannelParticipantsCount,
|
||
&i.ReplyChannelAdminsCount,
|
||
&i.ReplyChannelKickedCount,
|
||
&i.ReplyChannelBannedCount,
|
||
&i.ReplyChannelTopMessageID,
|
||
&i.ReplyChannelPinnedMessageID,
|
||
&i.ReplyChannelPts,
|
||
&i.ReplyChannelTtlPeriod,
|
||
&i.ReplyChannelDate,
|
||
&i.ReplyChannelDeleted,
|
||
); err != nil {
|
||
return nil, err
|
||
}
|
||
items = append(items, i)
|
||
}
|
||
if err := rows.Err(); err != nil {
|
||
return nil, err
|
||
}
|
||
return items, nil
|
||
}
|
||
|
||
const lockUserUpdateWatermark = `-- name: LockUserUpdateWatermark :one
|
||
SELECT contiguous_pts
|
||
FROM user_update_watermarks
|
||
WHERE user_id = $1
|
||
FOR UPDATE
|
||
`
|
||
|
||
func (q *Queries) LockUserUpdateWatermark(ctx context.Context, userID int64) (int32, error) {
|
||
row := q.db.QueryRow(ctx, lockUserUpdateWatermark, userID)
|
||
var contiguous_pts int32
|
||
err := row.Scan(&contiguous_pts)
|
||
return contiguous_pts, err
|
||
}
|
||
|
||
const markDispatchDelivered = `-- name: MarkDispatchDelivered :exec
|
||
DELETE FROM dispatch_outbox
|
||
WHERE target_user_id = $1
|
||
AND id = $2
|
||
`
|
||
|
||
type MarkDispatchDeliveredParams struct {
|
||
TargetUserID int64
|
||
ID int64
|
||
}
|
||
|
||
// 方案 A:投递成功即删除。outbox 是任务队列,delivered 行无保留价值
|
||
// (消息在 message_boxes、离线补偿在 user_update_events),删除让表维持「未完成任务」小稳态。
|
||
func (q *Queries) MarkDispatchDelivered(ctx context.Context, arg MarkDispatchDeliveredParams) error {
|
||
_, err := q.db.Exec(ctx, markDispatchDelivered, arg.TargetUserID, arg.ID)
|
||
return err
|
||
}
|
||
|
||
const markDispatchDeliveredBatch = `-- name: MarkDispatchDeliveredBatch :exec
|
||
DELETE FROM dispatch_outbox d
|
||
USING unnest($1::bigint[]) WITH ORDINALITY AS tu(target_user_id, ord)
|
||
JOIN unnest($2::bigint[]) WITH ORDINALITY AS di(id, ord) USING (ord)
|
||
WHERE d.target_user_id = tu.target_user_id
|
||
AND d.id = di.id
|
||
`
|
||
|
||
type MarkDispatchDeliveredBatchParams struct {
|
||
TargetUserIds []int64
|
||
Ids []int64
|
||
}
|
||
|
||
// 批量删除一批已投递的 (target_user_id, id);target_user_id 入 WHERE 保证分区裁剪。
|
||
func (q *Queries) MarkDispatchDeliveredBatch(ctx context.Context, arg MarkDispatchDeliveredBatchParams) error {
|
||
_, err := q.db.Exec(ctx, markDispatchDeliveredBatch, arg.TargetUserIds, arg.Ids)
|
||
return err
|
||
}
|
||
|
||
const markDispatchFailed = `-- name: MarkDispatchFailed :exec
|
||
UPDATE dispatch_outbox
|
||
SET
|
||
status = CASE WHEN attempts >= 5 THEN 'failed' ELSE 'pending' END,
|
||
next_attempt_at = CASE
|
||
WHEN attempts >= 5 THEN next_attempt_at
|
||
ELSE now() + make_interval(secs => LEAST(60, attempts * attempts))
|
||
END,
|
||
last_error = $3,
|
||
updated_at = now()
|
||
WHERE target_user_id = $1
|
||
AND id = $2
|
||
`
|
||
|
||
type MarkDispatchFailedParams struct {
|
||
TargetUserID int64
|
||
ID int64
|
||
LastError string
|
||
}
|
||
|
||
func (q *Queries) MarkDispatchFailed(ctx context.Context, arg MarkDispatchFailedParams) error {
|
||
_, err := q.db.Exec(ctx, markDispatchFailed, arg.TargetUserID, arg.ID, arg.LastError)
|
||
return err
|
||
}
|
||
|
||
const maxUserPts = `-- name: MaxUserPts :one
|
||
SELECT COALESCE(MAX(pts), 0)::int AS max_pts
|
||
FROM user_update_events
|
||
WHERE user_id = $1
|
||
`
|
||
|
||
func (q *Queries) MaxUserPts(ctx context.Context, userID int64) (int32, error) {
|
||
row := q.db.QueryRow(ctx, maxUserPts, userID)
|
||
var max_pts int32
|
||
err := row.Scan(&max_pts)
|
||
return max_pts, err
|
||
}
|
||
|
||
const nextUserPtsAfter = `-- name: NextUserPtsAfter :many
|
||
SELECT pts, pts_count
|
||
FROM user_update_events
|
||
WHERE user_id = $1
|
||
AND pts > $2
|
||
ORDER BY pts ASC
|
||
LIMIT $3
|
||
`
|
||
|
||
type NextUserPtsAfterParams struct {
|
||
UserID int64
|
||
Pts int32
|
||
LimitCount int32
|
||
}
|
||
|
||
type NextUserPtsAfterRow struct {
|
||
Pts int32
|
||
PtsCount int32
|
||
}
|
||
|
||
func (q *Queries) NextUserPtsAfter(ctx context.Context, arg NextUserPtsAfterParams) ([]NextUserPtsAfterRow, error) {
|
||
rows, err := q.db.Query(ctx, nextUserPtsAfter, arg.UserID, arg.Pts, arg.LimitCount)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
var items []NextUserPtsAfterRow
|
||
for rows.Next() {
|
||
var i NextUserPtsAfterRow
|
||
if err := rows.Scan(&i.Pts, &i.PtsCount); err != nil {
|
||
return nil, err
|
||
}
|
||
items = append(items, i)
|
||
}
|
||
if err := rows.Err(); err != nil {
|
||
return nil, err
|
||
}
|
||
return items, nil
|
||
}
|
||
|
||
const recentUserPts = `-- name: RecentUserPts :many
|
||
SELECT pts, pts_count
|
||
FROM user_update_events
|
||
WHERE user_id = $1
|
||
ORDER BY pts DESC
|
||
LIMIT $2
|
||
`
|
||
|
||
type RecentUserPtsParams struct {
|
||
UserID int64
|
||
WindowSize int32
|
||
}
|
||
|
||
type RecentUserPtsRow struct {
|
||
Pts int32
|
||
PtsCount int32
|
||
}
|
||
|
||
// 取某 user 最近的一段 pts(降序),供计算「最大连续已提交 pts」用。
|
||
// 只看顶部窗口:瞬时空洞只可能出现在最近在途事务区,窗口足够大即可覆盖其下方连续。
|
||
func (q *Queries) RecentUserPts(ctx context.Context, arg RecentUserPtsParams) ([]RecentUserPtsRow, error) {
|
||
rows, err := q.db.Query(ctx, recentUserPts, arg.UserID, arg.WindowSize)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
var items []RecentUserPtsRow
|
||
for rows.Next() {
|
||
var i RecentUserPtsRow
|
||
if err := rows.Scan(&i.Pts, &i.PtsCount); err != nil {
|
||
return nil, err
|
||
}
|
||
items = append(items, i)
|
||
}
|
||
if err := rows.Err(); err != nil {
|
||
return nil, err
|
||
}
|
||
return items, nil
|
||
}
|
||
|
||
const saveUserUpdateWatermark = `-- name: SaveUserUpdateWatermark :exec
|
||
INSERT INTO user_update_watermarks (user_id, contiguous_pts, updated_at)
|
||
VALUES ($1, $2, now())
|
||
ON CONFLICT (user_id) DO UPDATE SET
|
||
contiguous_pts = GREATEST(user_update_watermarks.contiguous_pts, EXCLUDED.contiguous_pts),
|
||
updated_at = now()
|
||
`
|
||
|
||
type SaveUserUpdateWatermarkParams struct {
|
||
UserID int64
|
||
ContiguousPts int32
|
||
}
|
||
|
||
func (q *Queries) SaveUserUpdateWatermark(ctx context.Context, arg SaveUserUpdateWatermarkParams) error {
|
||
_, err := q.db.Exec(ctx, saveUserUpdateWatermark, arg.UserID, arg.ContiguousPts)
|
||
return err
|
||
}
|