5 Commits

Author SHA1 Message Date
conatum d7ef852ebb fix: Correct notification subscription parameters 2026-07-22 13:10:40 +03:00
conatum 0326220880 fix: Correct subscription name 2026-07-22 13:08:57 +03:00
conatum 8791c37473 fix: Incorrect payload type in ChatNotification event 2026-07-22 13:07:37 +03:00
conatum 0c0040d3d2 (data): Add v1-v2 schema migration 2026-07-22 12:59:56 +03:00
conatum dd6fc013ee feat: Add channel notification event tracking.
Bumps schema version to 2.
Adds Profiles module version dep on schema init.
Add notice_events, notice_sub_events, and notice_resub_events tables.
Add corresponding event manager.
Add ChatNotificationWithPayload event payload extension.
Add ChannelNotification subscription.
2026-07-22 12:21:16 +03:00
5 changed files with 476 additions and 110 deletions
+166
View File
@@ -0,0 +1,166 @@
BEGIN;
ALTER TABLE
events
ADD COLUMN
event_payload JSONB;
ALTER TABLE
events
ADD COLUMN
parent_event_id INTEGER REFERENCES events (event_id) ON DELETE SET NULL ON UPDATE CASCADE;
CREATE TABLE notice_events(
event_id INTEGER PRIMARY KEY REFERENCES events (event_id),
event_type TEXT NOT NULL DEFAULT 'notice',
message_id TEXT NOT NULL,
system_message TEXT NOT NULL,
FOREIGN KEY (event_id, event_type) REFERENCES events (event_id, event_type)
);
CREATE TABLE notice_sub_events(
event_id INTEGER PRIMARY KEY REFERENCES events (event_id),
event_type TEXT NOT NULL DEFAULT 'notice_sub',
tier INTEGER NOT NULL,
is_prime BOOLEAN NOT NULL,
duration_months INTEGER NOT NULL,
FOREIGN KEY (event_id, event_type) REFERENCES events (event_id, event_type)
);
CREATE TABLE notice_resub_events(
event_id INTEGER PRIMARY KEY REFERENCES events (event_id),
event_type TEXT NOT NULL DEFAULT 'notice_resub',
tier INTEGER NOT NULL,
is_prime BOOLEAN,
is_gift BOOLEAN NOT NULL,
cumulative_months INTEGER NOT NULL,
duration_months INTEGER NOT NULL,
streak_months INTEGER,
gifter_is_anonymous BOOLEAN,
gifter_user_id TEXT,
FOREIGN KEY (event_id, event_type) REFERENCES events (event_id, event_type)
);
CREATE OR REPLACE VIEW event_info AS
SELECT
events.event_id AS event_id,
events.event_type AS event_type,
events.communityid AS event_communityid,
events.channel_id AS event_channel_id,
events.profileid AS user_profileid,
events.user_id AS user_id,
user_profiles.nickname AS user_nickname,
user_profiles.timezone AS user_timezone,
user_profiles.locale_hint AS user_localehint,
user_profiles.locale AS user_locale,
COALESCE(events.occurred_at, events.created_at) AS occurred_at,
events.created_at AS created_at,
follow_events.follower_count AS follower_count,
bits_events.bits AS bits_count,
bits_events.bits_type AS bits_type,
bits_events.message AS bits_message,
bits_events.powerup_type AS bits_powerup_type,
subscribe_events.tier AS subscription_tier,
subscribe_events.gifted AS subscription_gifted,
gift_events.tier AS gifted_tier,
gift_events.gifted_count AS gifted_count,
subscribe_message_events.tier AS submessage_tier,
subscribe_message_events.duration_months AS submessage_duration,
subscribe_message_events.cumulative_months AS submessage_cumulative,
subscribe_message_events.streak_months AS submessage_streak,
subscribe_message_events.message AS submessage_message,
cheer_events.amount AS cheer_amount,
cheer_events.message AS cheer_message,
redemption_add_events.redeem_id AS redeemed_redeemid,
redemption_add_events.redeem_title AS redeemed_title,
redemption_add_events.redeem_cost AS redeemed_cost,
redemption_add_events.redemption_id AS redeemed_redemptionid,
redemption_add_events.redemption_status AS redeemed_status,
redemption_add_events.message AS redeemed_message,
redemption_update_events.redeem_id AS redeemupdate_redeemid,
redemption_update_events.redeem_title AS redeemupdate_title,
redemption_update_events.redeem_cost AS redeemupdate_cost,
redemption_update_events.redemption_id AS redeemupdate_redemptionid,
redemption_update_events.redemption_status AS redeemupdate_status,
redemption_update_events.redeemed_at AS redeemupdate_redeemedat,
poll_end_events.poll_id AS pollend_pollid,
poll_end_events.poll_title AS pollend_title,
poll_end_events.poll_choices AS pollend_choices,
poll_end_events.poll_started AS pollend_started,
stream_online_events.stream_id AS online_streamid,
stream_online_events.stream_type AS online_streamtype,
channel_update_events.title AS update_title,
channel_update_events.language AS update_language,
channel_update_events.category_id AS update_catid,
channel_update_events.category_name AS update_catname,
message_events.message_id AS message_id,
message_events.message_type AS message_type,
message_events.content AS message_content,
message_events.source_channel_id AS message_sourceid,
raid_out_events.target_id AS raidout_targetid,
raid_out_events.target_name AS raidout_targetname,
raid_out_events.viewer_count AS raidout_viewers,
raid_in_events.source_id AS raidin_sourceid,
raid_in_events.source_name AS raidin_sourcename,
raid_in_events.viewer_count AS raidin_viewers,
notice_events.message_id AS notice_message_id,
notice_events.system_message AS notice_system_message,
notice_sub_events.tier AS notice_sub_tier,
notice_sub_events.is_prime AS notice_sub_is_prime,
notice_sub_events.duration_months AS notice_sub_duration_months,
events.event_payload AS event_payload,
events.parent_event_id AS parent_event_id,
notice_resub_events.tier AS notice_resub_tier,
notice_resub_events.is_prime AS notice_resub_is_prime,
notice_resub_events.is_gift AS notice_resub_is_gift,
notice_resub_events.cumulative_months AS notice_resub_cumulative_months,
notice_resub_events.duration_months AS notice_resub_duration_months,
notice_resub_events.streak_months AS notice_resub_streak_months,
notice_resub_events.gifter_is_anonymous AS notice_resub_gifter_is_anonymous,
notice_resub_events.gifter_user_id AS notice_resub_gifter_user_id
FROM events
LEFT JOIN user_profiles USING (profileid)
LEFT JOIN follow_events USING (event_id, event_type)
LEFT JOIN bits_events USING (event_id, event_type)
LEFT JOIN subscribe_events USING (event_id, event_type)
LEFT JOIN gift_events USING (event_id, event_type)
LEFT JOIN subscribe_message_events USING (event_id, event_type)
LEFT JOIN cheer_events USING (event_id, event_type)
LEFT JOIN redemption_add_events USING (event_id, event_type)
LEFT JOIN redemption_update_events USING (event_id, event_type)
LEFT JOIN poll_end_events USING (event_id, event_type)
LEFT JOIN stream_online_events USING (event_id, event_type)
LEFT JOIN stream_offline_events USING (event_id, event_type)
LEFT JOIN channel_update_events USING (event_id, event_type)
LEFT JOIN vip_add_events USING (event_id, event_type)
LEFT JOIN vip_remove_events USING (event_id, event_type)
LEFT JOIN message_events USING (event_id, event_type)
LEFT JOIN raid_out_events USING (event_id, event_type)
LEFT JOIN raid_in_events USING (event_id, event_type)
LEFT JOIN notice_events USING (event_id, event_type)
LEFT JOIN notice_sub_events USING (event_id, event_type)
LEFT JOIN notice_resub_events USING (event_id, event_type);
INSERT INTO version_history (component, from_version, to_version, author) VALUES ('EVENT_TRACKER', 1, 2, 'Version 1-2 migration');
COMMIT;
+63 -4
View File
@@ -1,8 +1,11 @@
BEGIN; BEGIN;
INSERT INTO version_history (component, from_version, to_version, author) VALUES ('EVENT_TRACKER', 0, 1, 'Initial Creation'); -- Version dependency checks
DO $$
ASSERT current_module_version('PROFILES') = 1, 'Dependency version mismatch: PROFILES';
$$ LANGUAGE plpgsql;
-- TODO: Assert profiles data version INSERT INTO version_history (component, from_version, to_version, author) VALUES ('EVENT_TRACKER', 0, 2, 'Initial Creation');
-- Twitch tracked event data {{{ -- Twitch tracked event data {{{
@@ -23,6 +26,8 @@ CREATE TABLE events(
profileid INTEGER REFERENCES user_profiles (profileid), profileid INTEGER REFERENCES user_profiles (profileid),
user_id TEXT, user_id TEXT,
occurred_at TIMESTAMPTZ, occurred_at TIMESTAMPTZ,
event_payload JSONB,
parent_event_id INTEGER REFERENCES events (event_id) ON DELETE SET NULL ON UPDATE CASCADE,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
UNIQUE (event_id, event_type) UNIQUE (event_id, event_type)
); );
@@ -180,6 +185,37 @@ CREATE TABLE raid_in_events(
FOREIGN KEY (event_id, event_type) REFERENCES events (event_id, event_type) FOREIGN KEY (event_id, event_type) REFERENCES events (event_id, event_type)
); );
CREATE TABLE notice_events(
event_id INTEGER PRIMARY KEY REFERENCES events (event_id),
event_type TEXT NOT NULL DEFAULT 'notice',
message_id TEXT NOT NULL,
system_message TEXT NOT NULL,
FOREIGN KEY (event_id, event_type) REFERENCES events (event_id, event_type)
);
CREATE TABLE notice_sub_events(
event_id INTEGER PRIMARY KEY REFERENCES events (event_id),
event_type TEXT NOT NULL DEFAULT 'notice_sub',
tier INTEGER NOT NULL,
is_prime BOOLEAN NOT NULL,
duration_months INTEGER NOT NULL,
FOREIGN KEY (event_id, event_type) REFERENCES events (event_id, event_type)
);
CREATE TABLE notice_resub_events(
event_id INTEGER PRIMARY KEY REFERENCES events (event_id),
event_type TEXT NOT NULL DEFAULT 'notice_resub',
tier INTEGER NOT NULL,
is_prime BOOLEAN,
is_gift BOOLEAN NOT NULL,
cumulative_months INTEGER NOT NULL,
duration_months INTEGER NOT NULL,
streak_months INTEGER,
gifter_is_anonymous BOOLEAN,
gifter_user_id TEXT,
FOREIGN KEY (event_id, event_type) REFERENCES events (event_id, event_type)
);
CREATE VIEW event_info AS CREATE VIEW event_info AS
SELECT SELECT
@@ -189,6 +225,7 @@ SELECT
events.channel_id AS event_channel_id, events.channel_id AS event_channel_id,
events.profileid AS user_profileid, events.profileid AS user_profileid,
events.user_id AS user_id, events.user_id AS user_id,
user_profiles.nickname AS user_nickname, user_profiles.nickname AS user_nickname,
user_profiles.timezone AS user_timezone, user_profiles.timezone AS user_timezone,
user_profiles.locale_hint AS user_localehint, user_profiles.locale_hint AS user_localehint,
@@ -256,7 +293,26 @@ SELECT
raid_in_events.source_id AS raidin_sourceid, raid_in_events.source_id AS raidin_sourceid,
raid_in_events.source_name AS raidin_sourcename, raid_in_events.source_name AS raidin_sourcename,
raid_in_events.viewer_count AS raidin_viewers raid_in_events.viewer_count AS raidin_viewers,
notice_events.message_id AS notice_message_id,
notice_events.system_message AS notice_system_message,
notice_sub_events.tier AS notice_sub_tier,
notice_sub_events.is_prime AS notice_sub_is_prime,
notice_sub_events.duration_months AS notice_sub_duration_months,
events.event_payload AS event_payload,
events.parent_event_id AS parent_event_id,
notice_resub_events.tier AS notice_resub_tier,
notice_resub_events.is_prime AS notice_resub_is_prime,
notice_resub_events.is_gift AS notice_resub_is_gift,
notice_resub_events.cumulative_months AS notice_resub_cumulative_months,
notice_resub_events.duration_months AS notice_resub_duration_months,
notice_resub_events.streak_months AS notice_resub_streak_months,
notice_resub_events.gifter_is_anonymous AS notice_resub_gifter_is_anonymous,
notice_resub_events.gifter_user_id AS notice_resub_gifter_user_id
FROM events FROM events
LEFT JOIN user_profiles USING (profileid) LEFT JOIN user_profiles USING (profileid)
@@ -276,7 +332,10 @@ LEFT JOIN vip_add_events USING (event_id, event_type)
LEFT JOIN vip_remove_events USING (event_id, event_type) LEFT JOIN vip_remove_events USING (event_id, event_type)
LEFT JOIN message_events USING (event_id, event_type) LEFT JOIN message_events USING (event_id, event_type)
LEFT JOIN raid_out_events USING (event_id, event_type) LEFT JOIN raid_out_events USING (event_id, event_type)
LEFT JOIN raid_in_events USING (event_id, event_type); LEFT JOIN raid_in_events USING (event_id, event_type)
LEFT JOIN notice_events USING (event_id, event_type)
LEFT JOIN notice_sub_events USING (event_id, event_type)
LEFT JOIN notice_resub_events USING (event_id, event_type);
-- }}} -- }}}
+206 -85
View File
@@ -1,8 +1,10 @@
from typing import Optional from typing import Optional
import random import random
import twitchio import twitchio
from twitchio import PartialUser, Scopes, eventsub from twitchio import PartialUser, Scopes, eventsub
from twitchio.ext import commands as cmds from twitchio.ext import commands as cmds
from psycopg.types.json import Jsonb
from botdata import BotChannel from botdata import BotChannel
from meta import Bot from meta import Bot
@@ -10,6 +12,7 @@ from utils.lib import utc_now
from . import logger from . import logger
from .data import EventData, TrackingChannel from .data import EventData, TrackingChannel
from .payloads import ChatNotificationWithPayload
class TrackerComponent(cmds.Component): class TrackerComponent(cmds.Component):
@@ -32,43 +35,49 @@ class TrackerComponent(cmds.Component):
# ----- Methods ----- # ----- Methods -----
async def start_tracking(self, channel: TrackingChannel): async def start_tracking(self, channel: TrackingChannel):
# TODO: Make sure that we aren't trying to make duplicate subscriptions here # TODO: Make sure that we aren't trying to make duplicate subscriptions here
logger.debug( logger.debug("Initialising event tracking for %s", channel.userid)
"Initialising event tracking for %s",
channel.userid
)
# Get associated auth scopes # Get associated auth scopes
rows = await self.bot.data.user_auth_scopes.select_where(userid=channel.userid) rows = await self.bot.data.user_auth_scopes.select_where(userid=channel.userid)
scopes = Scopes([row['scope'] for row in rows]) scopes = Scopes([row["scope"] for row in rows])
# Build subscription payloads based on available scopes # Build subscription payloads based on available scopes
subs = [] subs = []
usersubs = [] usersubs = []
subcls = [] subcls = []
if Scopes.channel_read_redemptions in scopes or Scopes.channel_manage_redemptions in scopes: if (
Scopes.channel_read_redemptions in scopes
or Scopes.channel_manage_redemptions in scopes
):
subcls.append(eventsub.ChannelPointsRedeemAddSubscription) subcls.append(eventsub.ChannelPointsRedeemAddSubscription)
subcls.append(eventsub.ChannelPointsRedeemUpdateSubscription) subcls.append(eventsub.ChannelPointsRedeemUpdateSubscription)
if Scopes.bits_read in scopes: if Scopes.bits_read in scopes:
subcls.append(eventsub.ChannelBitsUseSubscription) subcls.append(eventsub.ChannelBitsUseSubscription)
subcls.append(eventsub.ChannelCheerSubscription) subcls.append(eventsub.ChannelCheerSubscription)
if Scopes.channel_read_subscriptions in scopes: if Scopes.channel_read_subscriptions in scopes:
subcls.extend(( subcls.extend(
eventsub.ChannelSubscribeSubscription, (
eventsub.ChannelSubscribeMessageSubscription, eventsub.ChannelSubscribeSubscription,
eventsub.ChannelSubscriptionGiftSubscription, eventsub.ChannelSubscribeMessageSubscription,
)) eventsub.ChannelSubscriptionGiftSubscription,
)
)
if Scopes.channel_read_polls in scopes or Scopes.channel_manage_polls in scopes: if Scopes.channel_read_polls in scopes or Scopes.channel_manage_polls in scopes:
subcls.append(eventsub.ChannelPollEndSubscription) subcls.append(eventsub.ChannelPollEndSubscription)
if Scopes.channel_read_vips in scopes or Scopes.channel_manage_vips in scopes: if Scopes.channel_read_vips in scopes or Scopes.channel_manage_vips in scopes:
subcls.extend(( subcls.extend(
eventsub.ChannelVIPAddSubscription, (
eventsub.ChannelVIPRemoveSubscription, eventsub.ChannelVIPAddSubscription,
)) eventsub.ChannelVIPRemoveSubscription,
)
)
subcls.extend(( subcls.extend(
eventsub.StreamOnlineSubscription, (
eventsub.StreamOfflineSubscription, eventsub.StreamOnlineSubscription,
eventsub.ChannelUpdateSubscription, eventsub.StreamOfflineSubscription,
)) eventsub.ChannelUpdateSubscription,
)
)
for subbr in subcls: for subbr in subcls:
subs.append(subbr(broadcaster_user_id=channel.userid)) subs.append(subbr(broadcaster_user_id=channel.userid))
@@ -90,10 +99,18 @@ class TrackerComponent(cmds.Component):
) )
) )
subs.extend(( subs.extend(
eventsub.ChannelRaidSubscription(to_broadcaster_user_id=channel.userid), (
eventsub.ChannelRaidSubscription(from_broadcaster_user_id=channel.userid), eventsub.ChannelRaidSubscription(to_broadcaster_user_id=channel.userid),
)) eventsub.ChannelRaidSubscription(
from_broadcaster_user_id=channel.userid
),
eventsub.ChatNotificationSubscription(
broadcaster_user_id=channel.userid,
user_id=self.bot.bot_id,
),
)
)
responses = [] responses = []
for sub in subs: for sub in subs:
@@ -110,25 +127,33 @@ class TrackerComponent(cmds.Component):
if self.bot.using_webhooks: if self.bot.using_webhooks:
resp = await self.bot.subscribe_webhook(sub) resp = await self.bot.subscribe_webhook(sub)
else: else:
resp = await self.bot.subscribe_websocket(sub, token_for=channel.userid, as_bot=False) resp = await self.bot.subscribe_websocket(
sub, token_for=channel.userid, as_bot=False
)
responses.append(resp) responses.append(resp)
except Exception: except Exception:
logger.exception("Failed to subscribe to %s", str(sub)) logger.exception("Failed to subscribe to %s", str(sub))
logger.info("Finished tracker subscription to %s: %s", channel.userid, ', '.join(map(str, responses))) logger.info(
"Finished tracker subscription to %s: %s",
channel.userid,
", ".join(map(str, responses)),
)
# ----- Events ----- # ----- Events -----
@cmds.Component.listener() @cmds.Component.listener()
async def event_safe_channel_joined(self, payload: BotChannel): async def event_safe_channel_joined(self, payload: BotChannel):
# Check if the channel is tracked # Check if the channel is tracked
# If it is, call start_tracking # If it is, call start_tracking
tracked = await TrackingChannel.fetch(payload.userid) tracked = await TrackingChannel.fetch(payload.userid)
if tracked and tracked.joined: if tracked and tracked.joined:
await self.start_tracking(tracked) await self.start_tracking(tracked)
@cmds.Component.listener() @cmds.Component.listener()
async def event_custom_redemption_add(self, payload: twitchio.ChannelPointsRedemptionAdd): async def event_custom_redemption_add(
self, payload: twitchio.ChannelPointsRedemptionAdd
):
tracked = await TrackingChannel.fetch(payload.broadcaster.id) tracked = await TrackingChannel.fetch(payload.broadcaster.id)
if tracked and tracked.joined: if tracked and tracked.joined:
community = await self.bot.profiles.fetch_community(payload.broadcaster) community = await self.bot.profiles.fetch_community(payload.broadcaster)
@@ -137,7 +162,7 @@ class TrackerComponent(cmds.Component):
pid = profile.profileid pid = profile.profileid
event_row = await self.data.events.insert( event_row = await self.data.events.insert(
event_type='redemption_add', event_type="redemption_add",
communityid=cid, communityid=cid,
channel_id=payload.broadcaster.id, channel_id=payload.broadcaster.id,
profileid=pid, profileid=pid,
@@ -145,7 +170,7 @@ class TrackerComponent(cmds.Component):
occurred_at=payload.redeemed_at, occurred_at=payload.redeemed_at,
) )
detail_row = await self.data.redemption_add_events.insert( detail_row = await self.data.redemption_add_events.insert(
event_id=event_row['event_id'], event_id=event_row["event_id"],
redeem_id=payload.reward.id, redeem_id=payload.reward.id,
redeem_title=payload.reward.title, redeem_title=payload.reward.title,
redeem_cost=payload.reward.cost, redeem_cost=payload.reward.cost,
@@ -155,7 +180,9 @@ class TrackerComponent(cmds.Component):
) )
@cmds.Component.listener() @cmds.Component.listener()
async def event_custom_redemption_update(self, payload: twitchio.ChannelPointsRedemptionUpdate): async def event_custom_redemption_update(
self, payload: twitchio.ChannelPointsRedemptionUpdate
):
tracked = await TrackingChannel.fetch(payload.broadcaster.id) tracked = await TrackingChannel.fetch(payload.broadcaster.id)
if tracked and tracked.joined: if tracked and tracked.joined:
community = await self.bot.profiles.fetch_community(payload.broadcaster) community = await self.bot.profiles.fetch_community(payload.broadcaster)
@@ -164,20 +191,20 @@ class TrackerComponent(cmds.Component):
pid = profile.profileid pid = profile.profileid
event_row = await self.data.events.insert( event_row = await self.data.events.insert(
event_type='redemption_update', event_type="redemption_update",
communityid=cid, communityid=cid,
channel_id=payload.broadcaster.id, channel_id=payload.broadcaster.id,
profileid=pid, profileid=pid,
user_id=payload.user.id, user_id=payload.user.id,
) )
detail_row = await self.data.redemption_update_events.insert( detail_row = await self.data.redemption_update_events.insert(
event_id=event_row['event_id'], event_id=event_row["event_id"],
redeem_id=payload.reward.id, redeem_id=payload.reward.id,
redeem_title=payload.reward.title, redeem_title=payload.reward.title,
redeem_cost=payload.reward.cost, redeem_cost=payload.reward.cost,
redemption_id=payload.id, redemption_id=payload.id,
redemption_status=payload.status, redemption_status=payload.status,
redeemed_at=utc_now() redeemed_at=utc_now(),
) )
@cmds.Component.listener() @cmds.Component.listener()
@@ -189,19 +216,18 @@ class TrackerComponent(cmds.Component):
profile = await self.bot.profiles.fetch_profile(payload.user) profile = await self.bot.profiles.fetch_profile(payload.user)
pid = profile.profileid pid = profile.profileid
# Computer follower count # Computer follower count
followers = await payload.broadcaster.fetch_followers() followers = await payload.broadcaster.fetch_followers()
event_row = await self.data.events.insert( event_row = await self.data.events.insert(
event_type='follow', event_type="follow",
communityid=cid, communityid=cid,
channel_id=payload.broadcaster.id, channel_id=payload.broadcaster.id,
profileid=pid, profileid=pid,
user_id=payload.user.id, user_id=payload.user.id,
) )
detail_row = await self.data.follow_events.insert( detail_row = await self.data.follow_events.insert(
event_id=event_row['event_id'], event_id=event_row["event_id"], follower_count=followers.total
follower_count=followers.total
) )
@cmds.Component.listener() @cmds.Component.listener()
@@ -214,20 +240,20 @@ class TrackerComponent(cmds.Component):
pid = profile.profileid pid = profile.profileid
event_row = await self.data.events.insert( event_row = await self.data.events.insert(
event_type='bits', event_type="bits",
communityid=cid, communityid=cid,
channel_id=payload.broadcaster.id, channel_id=payload.broadcaster.id,
profileid=pid, profileid=pid,
user_id=payload.user.id, user_id=payload.user.id,
) )
detail_row = await self.data.bits_events.insert( detail_row = await self.data.bits_events.insert(
event_id=event_row['event_id'], event_id=event_row["event_id"],
bits=payload.bits, bits=payload.bits,
bits_type=payload.type, bits_type=payload.type,
message=payload.text, message=payload.text,
powerup_type=payload.power_up.type if payload.power_up else None powerup_type=payload.power_up.type if payload.power_up else None,
) )
self.bot.safe_dispatch('bits_use', payload=(event_row, detail_row, payload)) self.bot.safe_dispatch("bits_use", payload=(event_row, detail_row, payload))
@cmds.Component.listener() @cmds.Component.listener()
async def event_subscription(self, payload: twitchio.ChannelSubscribe): async def event_subscription(self, payload: twitchio.ChannelSubscribe):
@@ -239,18 +265,20 @@ class TrackerComponent(cmds.Component):
pid = profile.profileid pid = profile.profileid
event_row = await self.data.events.insert( event_row = await self.data.events.insert(
event_type='subscribe', event_type="subscribe",
communityid=cid, communityid=cid,
channel_id=payload.broadcaster.id, channel_id=payload.broadcaster.id,
profileid=pid, profileid=pid,
user_id=payload.user.id, user_id=payload.user.id,
) )
detail_row = await self.data.subscribe_events.insert( detail_row = await self.data.subscribe_events.insert(
event_id=event_row['event_id'], event_id=event_row["event_id"],
tier=int(payload.tier), tier=int(payload.tier),
gifted=payload.gift, gifted=payload.gift,
) )
self.bot.safe_dispatch('subscription', payload=(event_row, detail_row, payload)) self.bot.safe_dispatch(
"subscription", payload=(event_row, detail_row, payload)
)
@cmds.Component.listener() @cmds.Component.listener()
async def event_subscription_gift(self, payload: twitchio.ChannelSubscriptionGift): async def event_subscription_gift(self, payload: twitchio.ChannelSubscriptionGift):
@@ -265,21 +293,25 @@ class TrackerComponent(cmds.Component):
pid = None pid = None
event_row = await self.data.events.insert( event_row = await self.data.events.insert(
event_type='gift', event_type="gift",
communityid=cid, communityid=cid,
channel_id=payload.broadcaster.id, channel_id=payload.broadcaster.id,
profileid=pid, profileid=pid,
user_id=payload.user.id if payload.user else None, user_id=payload.user.id if payload.user else None,
) )
detail_row = await self.data.gift_events.insert( detail_row = await self.data.gift_events.insert(
event_id=event_row['event_id'], event_id=event_row["event_id"],
tier=int(payload.tier), tier=int(payload.tier),
gifted_count=payload.total, gifted_count=payload.total,
) )
self.bot.safe_dispatch('subscription_gift', payload=(event_row, detail_row, payload)) self.bot.safe_dispatch(
"subscription_gift", payload=(event_row, detail_row, payload)
)
@cmds.Component.listener() @cmds.Component.listener()
async def event_subscription_message(self, payload: twitchio.ChannelSubscriptionMessage): async def event_subscription_message(
self, payload: twitchio.ChannelSubscriptionMessage
):
tracked = await TrackingChannel.fetch(payload.broadcaster.id) tracked = await TrackingChannel.fetch(payload.broadcaster.id)
if tracked and tracked.joined: if tracked and tracked.joined:
community = await self.bot.profiles.fetch_community(payload.broadcaster) community = await self.bot.profiles.fetch_community(payload.broadcaster)
@@ -288,21 +320,23 @@ class TrackerComponent(cmds.Component):
pid = profile.profileid pid = profile.profileid
event_row = await self.data.events.insert( event_row = await self.data.events.insert(
event_type='subscribe_message', event_type="subscribe_message",
communityid=cid, communityid=cid,
channel_id=payload.broadcaster.id, channel_id=payload.broadcaster.id,
profileid=pid, profileid=pid,
user_id=payload.user.id, user_id=payload.user.id,
) )
detail_row = await self.data.subscribe_message_events.insert( detail_row = await self.data.subscribe_message_events.insert(
event_id=event_row['event_id'], event_id=event_row["event_id"],
tier=int(payload.tier), tier=int(payload.tier),
duration_months=payload.months, duration_months=payload.months,
cumulative_months=payload.cumulative_months, cumulative_months=payload.cumulative_months,
streak_months=payload.streak_months, streak_months=payload.streak_months,
message=payload.text, message=payload.text,
) )
self.bot.safe_dispatch('subscription_message', payload=(event_row, detail_row, payload)) self.bot.safe_dispatch(
"subscription_message", payload=(event_row, detail_row, payload)
)
@cmds.Component.listener() @cmds.Component.listener()
async def event_stream_online(self, payload: twitchio.StreamOnline): async def event_stream_online(self, payload: twitchio.StreamOnline):
@@ -312,13 +346,13 @@ class TrackerComponent(cmds.Component):
cid = community.communityid cid = community.communityid
event_row = await self.data.events.insert( event_row = await self.data.events.insert(
event_type='stream_online', event_type="stream_online",
communityid=cid, communityid=cid,
channel_id=payload.broadcaster.id, channel_id=payload.broadcaster.id,
occurred_at=payload.started_at, occurred_at=payload.started_at,
) )
detail_row = await self.data.stream_online_events.insert( detail_row = await self.data.stream_online_events.insert(
event_id=event_row['event_id'], event_id=event_row["event_id"],
stream_id=payload.id, stream_id=payload.id,
stream_type=payload.type, stream_type=payload.type,
) )
@@ -331,12 +365,12 @@ class TrackerComponent(cmds.Component):
cid = community.communityid cid = community.communityid
event_row = await self.data.events.insert( event_row = await self.data.events.insert(
event_type='stream_offline', event_type="stream_offline",
communityid=cid, communityid=cid,
channel_id=payload.broadcaster.id, channel_id=payload.broadcaster.id,
) )
detail_row = await self.data.stream_offline_events.insert( detail_row = await self.data.stream_offline_events.insert(
event_id=event_row['event_id'], event_id=event_row["event_id"],
) )
@cmds.Component.listener() @cmds.Component.listener()
@@ -352,42 +386,121 @@ class TrackerComponent(cmds.Component):
payload.viewer_count, payload.viewer_count,
) )
async def _event_raid_out(self, broadcaster: PartialUser, to_broadcaster: PartialUser, viewer_count: int): async def _event_raid_out(
self, broadcaster: PartialUser, to_broadcaster: PartialUser, viewer_count: int
):
tracked = await TrackingChannel.fetch(broadcaster.id) tracked = await TrackingChannel.fetch(broadcaster.id)
if tracked and tracked.joined: if tracked and tracked.joined:
community = await self.bot.profiles.fetch_community(broadcaster) community = await self.bot.profiles.fetch_community(broadcaster)
cid = community.communityid cid = community.communityid
event_row = await self.data.events.insert( event_row = await self.data.events.insert(
event_type='raidout', event_type="raidout",
communityid=cid, communityid=cid,
channel_id=broadcaster.id, channel_id=broadcaster.id,
) )
detail_row = await self.data.raid_out_events.insert( detail_row = await self.data.raid_out_events.insert(
event_id=event_row['event_id'], event_id=event_row["event_id"],
target_id=to_broadcaster.id, target_id=to_broadcaster.id,
target_name=to_broadcaster.name, target_name=to_broadcaster.name,
viewer_count=viewer_count viewer_count=viewer_count,
) )
async def _event_raid_in(self, broadcaster: PartialUser, from_broadcaster: PartialUser, viewer_count: int): async def _event_raid_in(
self, broadcaster: PartialUser, from_broadcaster: PartialUser, viewer_count: int
):
tracked = await TrackingChannel.fetch(broadcaster.id) tracked = await TrackingChannel.fetch(broadcaster.id)
if tracked and tracked.joined: if tracked and tracked.joined:
community = await self.bot.profiles.fetch_community(broadcaster) community = await self.bot.profiles.fetch_community(broadcaster)
cid = community.communityid cid = community.communityid
event_row = await self.data.events.insert( event_row = await self.data.events.insert(
event_type='raidin', event_type="raidin",
communityid=cid, communityid=cid,
channel_id=broadcaster.id, channel_id=broadcaster.id,
) )
detail_row = await self.data.raid_in_events.insert( detail_row = await self.data.raid_in_events.insert(
event_id=event_row['event_id'], event_id=event_row["event_id"],
source_id=from_broadcaster.id, source_id=from_broadcaster.id,
source_name=from_broadcaster.name, source_name=from_broadcaster.name,
viewer_count=viewer_count viewer_count=viewer_count,
) )
@cmds.Component.listener()
async def event_chat_notification(self, payload: ChatNotificationWithPayload):
tracked = await TrackingChannel.fetch(payload.broadcaster.id)
if tracked and tracked.joined:
community = await self.bot.profiles.fetch_community(payload.broadcaster)
cid = community.communityid
profile = await self.bot.profiles.fetch_profile(payload.chatter)
pid = profile.profileid
event_payload = Jsonb(payload._raw)
event_row = await self.data.events.insert(
event_type="notice",
communityid=cid,
channel_id=payload.broadcaster.id,
profileid=pid,
user_id=payload.chatter.id,
event_payload=event_payload,
)
parent_event_id = event_row["event_id"]
detail_row = await self.data.notice_events.insert(
event_id=parent_event_id,
message_id=payload.id,
system_message=payload.system_message,
)
self.bot.safe_dispatch(
"chat_notice", payload=(event_row, detail_row, payload)
)
# Dispatch sub-events
if payload.notice_type == "sub":
event_row = await self.data.events.insert(
event_type="notice_sub",
communityid=cid,
channel_id=payload.broadcaster.id,
profileid=pid,
user_id=payload.chatter.id,
parent_event_id=parent_event_id,
)
detail_row = await self.data.notice_sub_events.insert(
event_id=event_row["event_id"],
tier=int(payload.tier),
is_prime=payload.sub.prime,
duration_months=payload.sub.months,
)
self.bot.safe_dispatch(
"chat_notice_sub", payload=(event_row, detail_row, payload)
)
elif payload.notice_type == "resub":
event_row = await self.data.events.insert(
event_type="notice_resub",
communityid=cid,
channel_id=payload.broadcaster.id,
profileid=pid,
user_id=payload.chatter.id,
parent_event_id=parent_event_id,
)
detail_row = await self.data.notice_resub_events.insert(
event_id=event_row["event_id"],
message_id=payload.id,
tier=payload.resub.tier,
is_prime=payload.resub.prime,
is_gift=payload.resub.gift,
cumulative_months=payload.resub.cumulative_months,
duration_months=payload.resub.months,
streak_months=payload.resub.streak_months,
gifter_is_anonymous=payload.resub.anonymous,
gifter_user_id=(
payload.resub.gifter.id if payload.resub.gifter else None
),
)
self.bot.safe_dispatch(
"chat_notice_resub", payload=(event_row, detail_row, payload)
)
@cmds.Component.listener() @cmds.Component.listener()
async def event_message(self, payload: twitchio.ChatMessage): async def event_message(self, payload: twitchio.ChatMessage):
tracked = await TrackingChannel.fetch(payload.broadcaster.id) tracked = await TrackingChannel.fetch(payload.broadcaster.id)
@@ -398,64 +511,72 @@ class TrackerComponent(cmds.Component):
pid = profile.profileid pid = profile.profileid
event_row = await self.data.events.insert( event_row = await self.data.events.insert(
event_type='message', event_type="message",
communityid=cid, communityid=cid,
channel_id=payload.broadcaster.id, channel_id=payload.broadcaster.id,
profileid=pid, profileid=pid,
user_id=payload.chatter.id, user_id=payload.chatter.id,
) )
detail_row = await self.data.message_events.insert( detail_row = await self.data.message_events.insert(
event_id=event_row['event_id'], event_id=event_row["event_id"],
message_id=payload.id, message_id=payload.id,
message_type=payload.type, message_type=payload.type,
content=payload.text, content=payload.text,
source_channel_id=payload.source_id source_channel_id=payload.source_id,
) )
# ----- Commands ----- # ----- Commands -----
@cmds.command(name='starttracking') @cmds.command(name="starttracking")
async def cmd_starttracking(self, ctx: cmds.Context): async def cmd_starttracking(self, ctx: cmds.Context):
if ctx.broadcaster: if ctx.broadcaster:
tracking = await TrackingChannel.fetch_or_create(ctx.channel.id, joined=True) tracking = await TrackingChannel.fetch_or_create(
ctx.channel.id, joined=True
)
if not tracking.joined: if not tracking.joined:
await tracking.update(joined=True) await tracking.update(joined=True)
rows = await self.bot.data.user_auth_scopes.select_where(userid=ctx.channel.id) rows = await self.bot.data.user_auth_scopes.select_where(
scopes = Scopes([row['scope'] for row in rows]) userid=ctx.channel.id
)
scopes = Scopes([row["scope"] for row in rows])
url = self.bot.get_auth_url( url = self.bot.get_auth_url(
Scopes({ Scopes(
Scopes.channel_read_subscriptions, {
Scopes.channel_read_redemptions, Scopes.channel_read_subscriptions,
Scopes.bits_read, Scopes.channel_read_redemptions,
Scopes.channel_read_polls, Scopes.bits_read,
Scopes.channel_read_vips, Scopes.channel_read_polls,
Scopes.moderator_read_followers, Scopes.channel_read_vips,
*scopes Scopes.moderator_read_followers,
}) *scopes,
}
)
)
await ctx.reply(
f"Tracking enabled! Please authorise me to track events in this channel: {url}"
) )
await ctx.reply(f"Tracking enabled! Please authorise me to track events in this channel: {url}")
else: else:
await ctx.reply("Only the broadcaster can enable tracking.") await ctx.reply("Only the broadcaster can enable tracking.")
@cmds.command(name='stoptracking') @cmds.command(name="stoptracking")
async def cmd_stoptracking(self, ctx: cmds.Context): async def cmd_stoptracking(self, ctx: cmds.Context):
if ctx.broadcaster: if ctx.broadcaster:
tracking = await TrackingChannel.fetch(ctx.channel.id) tracking = await TrackingChannel.fetch(ctx.channel.id)
if tracking and tracking.joined: if tracking and tracking.joined:
await tracking.update(joined=False) await tracking.update(joined=False)
# TODO: Actually disable the subscriptions instead of just on the next restart # TODO: Actually disable the subscriptions instead of just on the next restart
# This is tricky because some of the subscriptions may have been requested by other modules # This is tricky because some of the subscriptions may have been requested by other modules
# Requires keeping track of the source of subscriptions, and having a central manager disable them when no-one is listening anymore. # Requires keeping track of the source of subscriptions, and having a central manager disable them when no-one is listening anymore.
pass pass
await ctx.reply("Event tracking has been disabled.") await ctx.reply("Event tracking has been disabled.")
else: else:
await ctx.reply("Event tracking is not enabled!") await ctx.reply("Event tracking is not enabled!")
else: else:
await ctx.reply("Only the broadcaster can enable tracking.") await ctx.reply("Only the broadcaster can enable tracking.")
@cmds.command(name='join') @cmds.command(name="join")
async def cmd_join(self, ctx: cmds.Context): async def cmd_join(self, ctx: cmds.Context):
url = self.bot.get_auth_url() url = self.bot.get_auth_url()
await ctx.reply(f"Invite me to your channel with: {url}") await ctx.reply(f"Invite me to your channel with: {url}")
+25 -21
View File
@@ -3,39 +3,43 @@ from data.columns import String, Timestamp, Integer, Bool
class TrackingChannel(RowModel): class TrackingChannel(RowModel):
_tablename_ = 'tracking_channels' _tablename_ = "tracking_channels"
_cache_ = {} _cache_ = {}
userid = String(primary=True) userid = String(primary=True)
joined = Bool joined = Bool
joined_at = Timestamp() joined_at = Timestamp()
_timestamp = Timestamp() _timestamp = Timestamp()
class EventData(Registry): class EventData(Registry):
VERSION = ('EVENT_TRACKER', 1) VERSION = ("EVENT_TRACKER", 2)
tracking_channels = TrackingChannel.table tracking_channels = TrackingChannel.table
events = Table('events') events = Table("events")
follow_events = Table('follow_events') follow_events = Table("follow_events")
bits_events = Table('bits_events') bits_events = Table("bits_events")
subscribe_events = Table('subscribe_events') subscribe_events = Table("subscribe_events")
gift_events = Table('gift_events') gift_events = Table("gift_events")
subscribe_message_events = Table('subscribe_message_events') subscribe_message_events = Table("subscribe_message_events")
cheer_events = Table('cheer_events') cheer_events = Table("cheer_events")
redemption_add_events = Table('redemption_add_events') redemption_add_events = Table("redemption_add_events")
redemption_update_events = Table('redemption_update_events') redemption_update_events = Table("redemption_update_events")
poll_end_events = Table('poll_end_events') poll_end_events = Table("poll_end_events")
stream_online_events = Table('stream_online_events') stream_online_events = Table("stream_online_events")
stream_offline_events = Table('stream_offline_events') stream_offline_events = Table("stream_offline_events")
channel_update_events = Table('channel_update_events') channel_update_events = Table("channel_update_events")
vip_add_events = Table('vip_add_events') vip_add_events = Table("vip_add_events")
vip_remove_events = Table('vip_remove_events') vip_remove_events = Table("vip_remove_events")
raid_out_events = Table('raid_out_events') raid_out_events = Table("raid_out_events")
raid_in_events = Table('raid_in_events') raid_in_events = Table("raid_in_events")
message_events = Table('message_events') message_events = Table("message_events")
notice_events = Table("notice_events")
notice_sub_events = Table("notice_sub_events")
notice_resub_events = Table("notice_resub_events")
+16
View File
@@ -0,0 +1,16 @@
import twitchio
from twitchio.types_.eventsub import *
class ChatNotificationWithPayload(twitchio.ChatNotification):
"""
Extends twitchio.ChatNotification to include a _raw field with the original data.
The subclass registry in twitchio.BaseEvent should automatically register this as the correct event handler.
"""
__slots__ = ("_raw",)
def __init__(self, payload: ChannelChatNotificationEvent, *args, **kwargs):
super().__init__(payload, *args, **kwargs)
self._raw = payload