From 045e0f6db4ef86be67675b988aef8383d4d5eeeb Mon Sep 17 00:00:00 2001 From: A Date: Wed, 22 Jul 2026 21:44:33 +0800 Subject: [PATCH] fix: sync StarGift prepaid upgrade message refs --- ...35_star_gift_prepaid_message_refs.down.sql | 4 + ...0135_star_gift_prepaid_message_refs.up.sql | 304 ++++++++++++++ internal/store/postgres/star_gift.go | 3 +- .../store/postgres/star_gift_entitlements.go | 8 +- .../star_gift_lifecycle_integration_test.go | 189 ++++++++- ...ft_lifecycle_migration_integration_test.go | 4 +- internal/store/postgres/star_gift_upgrade.go | 377 ++++++++++++------ .../postgres/star_gift_user_message_ref.go | 13 +- 8 files changed, 754 insertions(+), 148 deletions(-) create mode 100644 deploy/migrations/0135_star_gift_prepaid_message_refs.down.sql create mode 100644 deploy/migrations/0135_star_gift_prepaid_message_refs.up.sql diff --git a/deploy/migrations/0135_star_gift_prepaid_message_refs.down.sql b/deploy/migrations/0135_star_gift_prepaid_message_refs.down.sql new file mode 100644 index 00000000..735e157c --- /dev/null +++ b/deploy/migrations/0135_star_gift_prepaid_message_refs.down.sql @@ -0,0 +1,4 @@ +-- The up migration registers protocol identities and emits durable per-user +-- edit_message events. Removing aliases, reverting snapshots or rewinding pts +-- would invalidate messages already consumed by clients and create holes in +-- updates.getDifference, so rollback intentionally preserves the repair. diff --git a/deploy/migrations/0135_star_gift_prepaid_message_refs.up.sql b/deploy/migrations/0135_star_gift_prepaid_message_refs.up.sql new file mode 100644 index 00000000..82cb120c --- /dev/null +++ b/deploy/migrations/0135_star_gift_prepaid_message_refs.up.sql @@ -0,0 +1,304 @@ +-- A separate prepaid-upgrade service message is another owner-local entry to +-- the same saved-gift aggregate. Earlier writes persisted gift_msg_id in the +-- receiver projection but did not register that message id, so clients that +-- submitted the visible card id received STARGIFT_INVALID. If the gift was +-- upgraded through the original id, the prepaid card also remained actionable. +-- +-- Repair aliases and already-upgraded projections atomically. Durable edit +-- events make history, online delivery and updates.getDifference converge on +-- the same non-actionable snapshot. Invalid persisted shapes fail the migration +-- instead of being normalized by a read path. + +LOCK TABLE public.peer_star_gifts, public.star_gift_user_message_refs, + public.message_boxes, public.private_messages IN SHARE ROW EXCLUSIVE MODE; + +DO $$ +BEGIN + IF EXISTS ( + SELECT 1 + FROM public.message_boxes box + WHERE NOT box.deleted + AND box.media #>> '{service_action,kind}' = 'star_gift' + AND box.media #>> '{service_action,star_gift,prepaid_upgrade}' = 'true' + AND box.media #>> '{service_action,star_gift,upgrade_separate}' = 'true' + AND ( + jsonb_typeof(box.media #> '{service_action,star_gift,gift_id}') IS DISTINCT FROM 'number' + OR COALESCE(box.media #>> '{service_action,star_gift,gift_id}', '') !~ '^[0-9]+$' + OR (box.media #>> '{service_action,star_gift,gift_id}')::numeric <= 0 + OR (box.media #>> '{service_action,star_gift,gift_id}')::numeric > 9223372036854775807 + ) + ) THEN + RAISE EXCEPTION 'separate prepaid star gift message has malformed gift_id'; + END IF; + + -- gift_msg_id is receiver-only, so absence is valid on the payer box. If + -- present it must be a positive protocol int32 message id. + IF EXISTS ( + SELECT 1 + FROM public.message_boxes box + WHERE NOT box.deleted + AND box.media #>> '{service_action,kind}' = 'star_gift' + AND box.media #>> '{service_action,star_gift,prepaid_upgrade}' = 'true' + AND box.media #>> '{service_action,star_gift,upgrade_separate}' = 'true' + AND box.media #> '{service_action,star_gift,gift_msg_id}' IS NOT NULL + AND ( + jsonb_typeof(box.media #> '{service_action,star_gift,gift_msg_id}') <> 'number' + OR COALESCE(box.media #>> '{service_action,star_gift,gift_msg_id}', '') !~ '^[0-9]+$' + OR (box.media #>> '{service_action,star_gift,gift_msg_id}')::numeric <= 0 + OR (box.media #>> '{service_action,star_gift,gift_msg_id}')::numeric > 2147483647 + ) + ) THEN + RAISE EXCEPTION 'separate prepaid star gift message has malformed gift_msg_id'; + END IF; +END +$$; + +CREATE TEMP TABLE star_gift_prepaid_message_aliases ON COMMIT DROP AS +SELECT DISTINCT owner_box.owner_user_id, + owner_box.box_id, + gift.id AS saved_gift_id, + owner_box.message_sender_id, + owner_box.private_message_id +FROM public.message_boxes owner_box +JOIN public.peer_star_gifts gift + ON gift.owner_peer_type = 'user' + AND gift.owner_peer_id = owner_box.owner_user_id + AND gift.lifecycle_status = 'active' + AND gift.msg_id = (owner_box.media #>> '{service_action,star_gift,gift_msg_id}')::integer + AND gift.gift_id = (owner_box.media #>> '{service_action,star_gift,gift_id}')::bigint +WHERE NOT owner_box.deleted + AND owner_box.media #>> '{service_action,kind}' = 'star_gift' + AND owner_box.media #>> '{service_action,star_gift,prepaid_upgrade}' = 'true' + AND owner_box.media #>> '{service_action,star_gift,upgrade_separate}' = 'true' + AND owner_box.media #> '{service_action,star_gift,gift_msg_id}' IS NOT NULL; + +DO $$ +BEGIN + IF EXISTS ( + SELECT 1 + FROM star_gift_prepaid_message_aliases + GROUP BY owner_user_id, box_id + HAVING COUNT(DISTINCT saved_gift_id) <> 1 + ) THEN + RAISE EXCEPTION 'separate prepaid star gift message resolves to multiple aggregates'; + END IF; + + IF EXISTS ( + SELECT 1 + FROM star_gift_prepaid_message_aliases alias + JOIN public.star_gift_user_message_refs ref + ON ref.owner_user_id = alias.owner_user_id + AND ref.msg_id = alias.box_id + WHERE ref.saved_gift_id <> alias.saved_gift_id + ) THEN + RAISE EXCEPTION 'separate prepaid star gift message collides with another aggregate'; + END IF; + + -- Both boxes of the logical private message must retain the same prepayment + -- identity. The receiver-only gift_msg_id may differ by design. + IF EXISTS ( + SELECT 1 + FROM star_gift_prepaid_message_aliases alias + JOIN public.peer_star_gifts gift ON gift.id = alias.saved_gift_id + JOIN public.message_boxes visible_box + ON visible_box.message_sender_id = alias.message_sender_id + AND visible_box.private_message_id = alias.private_message_id + AND NOT visible_box.deleted + WHERE visible_box.media #>> '{service_action,kind}' IS DISTINCT FROM 'star_gift' + OR visible_box.media #>> '{service_action,star_gift,prepaid_upgrade}' IS DISTINCT FROM 'true' + OR visible_box.media #>> '{service_action,star_gift,upgrade_separate}' IS DISTINCT FROM 'true' + OR visible_box.media #>> '{service_action,star_gift,gift_id}' IS DISTINCT FROM gift.gift_id::text + ) THEN + RAISE EXCEPTION 'separate prepaid star gift private projections disagree'; + END IF; +END +$$; + +CREATE UNIQUE INDEX star_gift_prepaid_message_aliases_owner_msg_idx + ON star_gift_prepaid_message_aliases(owner_user_id, box_id); + +INSERT INTO public.star_gift_user_message_refs(owner_user_id, msg_id, saved_gift_id) +SELECT owner_user_id, box_id, saved_gift_id +FROM star_gift_prepaid_message_aliases +ON CONFLICT (owner_user_id, msg_id) DO UPDATE +SET saved_gift_id = EXCLUDED.saved_gift_id +WHERE star_gift_user_message_refs.saved_gift_id = EXCLUDED.saved_gift_id; + +COMMENT ON TABLE public.star_gift_user_message_refs IS + 'Owner-local service-message aliases (unique outputs and separate prepaid-upgrade notifications) for one saved gift aggregate.'; + +CREATE TEMP TABLE star_gift_prepaid_message_repairs ( + owner_user_id bigint NOT NULL, + box_id integer NOT NULL, + peer_type text NOT NULL, + peer_id bigint NOT NULL, + message_sender_id bigint NOT NULL, + private_message_id bigint NOT NULL, + repaired_media jsonb NOT NULL, + PRIMARY KEY (owner_user_id, box_id) +) ON COMMIT DROP; + +-- Upgrade every visible copy of an already-consumed prepayment. A viewer gets +-- upgrade_msg_id only when that same viewer owns a box for the emitted unique +-- action. This covers the original sender while avoiding an owner-local link +-- on an unrelated third-party payer's card. +INSERT INTO star_gift_prepaid_message_repairs( + owner_user_id, box_id, peer_type, peer_id, + message_sender_id, private_message_id, repaired_media +) +SELECT visible_box.owner_user_id, + visible_box.box_id, + visible_box.peer_type, + visible_box.peer_id, + visible_box.message_sender_id, + visible_box.private_message_id, + CASE + WHEN unique_box.box_id IS NULL THEN + visible_box.media + #- '{service_action,star_gift,can_upgrade}' + #- '{service_action,star_gift,prepaid_upgrade_hash}' + #- '{service_action,star_gift,upgrade_msg_id}' + ELSE jsonb_set( + visible_box.media + #- '{service_action,star_gift,can_upgrade}' + #- '{service_action,star_gift,prepaid_upgrade_hash}', + '{service_action,star_gift,upgrade_msg_id}', + to_jsonb(unique_box.box_id::bigint), + true + ) + END +FROM star_gift_prepaid_message_aliases alias +JOIN public.peer_star_gifts gift + ON gift.id = alias.saved_gift_id + AND gift.lifecycle_status = 'active' + AND gift.unique_gift_id IS NOT NULL + AND gift.upgrade_msg_id > 0 +JOIN public.message_boxes owner_unique_box + ON owner_unique_box.owner_user_id = gift.owner_peer_id + AND owner_unique_box.box_id = gift.upgrade_msg_id + AND NOT owner_unique_box.deleted + AND owner_unique_box.media #>> '{service_action,kind}' = 'star_gift_unique' + AND owner_unique_box.media #>> '{service_action,star_gift_unique,gift,ID}' = gift.unique_gift_id::text +JOIN public.message_boxes visible_box + ON visible_box.message_sender_id = alias.message_sender_id + AND visible_box.private_message_id = alias.private_message_id + AND NOT visible_box.deleted +LEFT JOIN public.message_boxes unique_box + ON unique_box.owner_user_id = visible_box.owner_user_id + AND unique_box.message_sender_id = owner_unique_box.message_sender_id + AND unique_box.private_message_id = owner_unique_box.private_message_id + AND NOT unique_box.deleted + AND unique_box.media #>> '{service_action,kind}' = 'star_gift_unique' + AND unique_box.media #>> '{service_action,star_gift_unique,gift,ID}' = gift.unique_gift_id::text; + +DO $$ +DECLARE + repair_row record; + next_pts integer; + event_date integer := EXTRACT(EPOCH FROM clock_timestamp())::integer; +BEGIN + IF EXISTS ( + SELECT 1 + FROM star_gift_prepaid_message_aliases alias + JOIN public.peer_star_gifts gift + ON gift.id = alias.saved_gift_id + AND gift.lifecycle_status = 'active' + AND gift.unique_gift_id IS NOT NULL + WHERE NOT EXISTS ( + SELECT 1 + FROM star_gift_prepaid_message_repairs target_repair + WHERE target_repair.owner_user_id = alias.owner_user_id + AND target_repair.box_id = alias.box_id + ) + ) THEN + RAISE EXCEPTION 'upgraded star gift is missing its prepaid message repair'; + END IF; + + FOR repair_row IN + SELECT owner_user_id, box_id, peer_type, peer_id, repaired_media + FROM star_gift_prepaid_message_repairs + ORDER BY owner_user_id, box_id + LOOP + INSERT INTO public.user_update_watermarks(user_id, contiguous_pts) + VALUES(repair_row.owner_user_id, 0) + ON CONFLICT(user_id) DO NOTHING; + + UPDATE public.user_update_watermarks + SET contiguous_pts = contiguous_pts + 1, + updated_at = now() + WHERE user_id = repair_row.owner_user_id + RETURNING contiguous_pts INTO next_pts; + + UPDATE public.message_boxes + SET media = repair_row.repaired_media, + pts = next_pts + WHERE owner_user_id = repair_row.owner_user_id + AND box_id = repair_row.box_id + AND NOT deleted; + + INSERT INTO public.user_update_events( + user_id, pts, pts_count, date, event_type, + message_box_id, peer_type, peer_id + ) VALUES ( + repair_row.owner_user_id, next_pts, 1, event_date, 'edit_message', + repair_row.box_id, repair_row.peer_type, repair_row.peer_id + ); + + INSERT INTO public.dispatch_outbox( + target_user_id, pts, event_type, + exclude_auth_key_id, exclude_session_id + ) VALUES(repair_row.owner_user_id, next_pts, 'edit_message', 0, 0); + END LOOP; +END +$$; + +-- private_messages is a shared logical envelope and cannot retain either +-- participant's box-local gift_msg_id or upgrade_msg_id. +WITH shared_repairs AS ( + SELECT DISTINCT ON (repair.message_sender_id, repair.private_message_id) + repair.message_sender_id, + repair.private_message_id, + repair.repaired_media + #- '{service_action,star_gift,saved_id}' + #- '{service_action,star_gift,gift_msg_id}' + #- '{service_action,star_gift,upgrade_msg_id}' AS shared_media + FROM star_gift_prepaid_message_repairs repair + ORDER BY repair.message_sender_id, + repair.private_message_id, + (repair.owner_user_id = repair.message_sender_id) DESC, + repair.owner_user_id +) +UPDATE public.private_messages private_message +SET media = repair.shared_media +FROM shared_repairs repair +WHERE private_message.sender_user_id = repair.message_sender_id + AND private_message.id = repair.private_message_id; + +DO $$ +BEGIN + IF EXISTS ( + SELECT 1 + FROM star_gift_prepaid_message_aliases alias + LEFT JOIN public.star_gift_user_message_refs ref + ON ref.owner_user_id = alias.owner_user_id + AND ref.msg_id = alias.box_id + AND ref.saved_gift_id = alias.saved_gift_id + WHERE ref.saved_gift_id IS NULL + ) THEN + RAISE EXCEPTION 'separate prepaid star gift alias repair did not converge'; + END IF; + + IF EXISTS ( + SELECT 1 + FROM star_gift_prepaid_message_repairs repair + JOIN public.message_boxes box + ON box.owner_user_id = repair.owner_user_id + AND box.box_id = repair.box_id + WHERE box.media IS DISTINCT FROM repair.repaired_media + OR box.media #> '{service_action,star_gift,can_upgrade}' IS NOT NULL + OR box.media #> '{service_action,star_gift,prepaid_upgrade_hash}' IS NOT NULL + ) THEN + RAISE EXCEPTION 'upgraded prepaid star gift projection repair did not converge'; + END IF; +END +$$; diff --git a/internal/store/postgres/star_gift.go b/internal/store/postgres/star_gift.go index c488f6cd..87e69b61 100644 --- a/internal/store/postgres/star_gift.go +++ b/internal/store/postgres/star_gift.go @@ -655,7 +655,8 @@ LEFT JOIN unique_star_gifts u ON u.id=p.unique_gift_id CROSS JOIN LATERAL ( SELECT p.msg_id::bigint AS msg_id UNION ALL - SELECT r.msg_id::bigint FROM star_gift_user_message_refs r WHERE r.saved_gift_id=p.id + SELECT r.msg_id::bigint FROM star_gift_user_message_refs r + WHERE r.saved_gift_id=p.id AND r.owner_user_id=p.owner_peer_id ) ref WHERE p.owner_peer_type=$1 AND p.owner_peer_id=$2 AND p.lifecycle_status='active' AND (ref.msg_id=ANY($3::bigint[]) diff --git a/internal/store/postgres/star_gift_entitlements.go b/internal/store/postgres/star_gift_entitlements.go index dd261b1a..3d951201 100644 --- a/internal/store/postgres/star_gift_entitlements.go +++ b/internal/store/postgres/star_gift_entitlements.go @@ -140,8 +140,12 @@ VALUES($1,$2,$3,$4,$5,$6,$7)`, req.PayerUserID, req.CommandKey, locked.ID, req.F } return projectPrivateStarGiftSourceRef(ctx, tx, messageReq, result.Saved.Owner.ID, result.Saved.MsgID) }, after: func(ctx context.Context, tx pgx.Tx, sent domain.SendPrivateTextResult) error { - if req.Owner.Type != domain.PeerTypeChannel { - return nil + if req.Owner.Type == domain.PeerTypeUser { + ownerMessageID := sent.RecipientMessage.ID + if sent.SenderMessage.OwnerUserID == req.Owner.ID { + ownerMessageID = sent.SenderMessage.ID + } + return registerUserStarGiftMessageRef(ctx, tx, req.Owner.ID, ownerMessageID, result.Saved.ID, 0) } action := messageReq.Media.ServiceAction.StarGift return NewChannelStore(tx).appendStarGiftAdminLogTx(ctx, tx, req.Owner.ID, req.PayerUserID, diff --git a/internal/store/postgres/star_gift_lifecycle_integration_test.go b/internal/store/postgres/star_gift_lifecycle_integration_test.go index d91cb45f..62437833 100644 --- a/internal/store/postgres/star_gift_lifecycle_integration_test.go +++ b/internal/store/postgres/star_gift_lifecycle_integration_test.go @@ -7,6 +7,9 @@ import ( "testing" "time" + "github.com/jackc/pgx/v5/pgxpool" + + "telesrv/deploy" "telesrv/internal/domain" ) @@ -21,10 +24,11 @@ func TestStarGiftLifecycleAggregatePostgres(t *testing.T) { offerBuyer := createTestUser(t, ctx, users, "+1881"+suffix+"03", "OfferBuyer", "") resaleBuyer := createTestUser(t, ctx, users, "+1881"+suffix+"04", "ResaleBuyer", "") loser := createTestUser(t, ctx, users, "+1881"+suffix+"05", "AuctionLoser", "") + prepayPayer := createTestUser(t, ctx, users, "+1881"+suffix+"06", "PrepayPayer", "") ownerPeer := domain.Peer{Type: domain.PeerTypeUser, ID: owner.ID} stars := NewStarsStore(pool) - for _, user := range []domain.User{buyer, owner, offerBuyer, resaleBuyer, loser} { + for _, user := range []domain.User{buyer, owner, offerBuyer, resaleBuyer, loser, prepayPayer} { if _, _, err := stars.EnsureGrant(ctx, user.ID, 10000, now); err != nil { t.Fatalf("grant stars to %d: %v", user.ID, err) } @@ -100,10 +104,10 @@ func TestStarGiftLifecycleAggregatePostgres(t *testing.T) { t.Fatalf("prepaid target = %+v price %d err %v", target, price, err) } prepaid, err := lifecycle.PrepayStarGiftUpgrade(ctx, domain.StarGiftPrepaidUpgradeRequest{ - PayerUserID: buyer.ID, Owner: ownerPeer, Hash: purchased.Saved.PrepaidUpgradeHash, + PayerUserID: prepayPayer.ID, Owner: ownerPeer, Hash: purchased.Saved.PrepaidUpgradeHash, ChargeStars: 100, FormID: 11002, CommandKey: "prepay-" + suffix, Date: now + 1, }) - if err != nil || prepaid.Saved.PrepaidUpgradeStars != 100 || prepaid.Saved.PrepaidUpgradeHash != "" || prepaid.Balance.Balance != 9850 { + if err != nil || prepaid.Saved.PrepaidUpgradeStars != 100 || prepaid.Saved.PrepaidUpgradeHash != "" || prepaid.Balance.Balance != 9900 { t.Fatalf("prepay upgrade = %+v err %v", prepaid, err) } prepaySenderAction := prepaid.Send.SenderMessage.Media.ServiceAction.StarGift @@ -114,7 +118,7 @@ func TestStarGiftLifecycleAggregatePostgres(t *testing.T) { t.Fatalf("prepay gift_msg_id is not owner-only: sender=%+v owner=%+v purchase=%+v", prepaySenderAction, prepayOwnerAction, purchased.Send) } - prepaySenderDifference, err := NewUpdateEventStore(pool).ListAfter(ctx, buyer.ID, prepaid.Send.SenderMessage.Pts-1, 1) + prepaySenderDifference, err := NewUpdateEventStore(pool).ListAfter(ctx, prepayPayer.ID, prepaid.Send.SenderMessage.Pts-1, 1) if err != nil || len(prepaySenderDifference) != 1 || prepaySenderDifference[0].Message.Media == nil || prepaySenderDifference[0].Message.Media.ServiceAction == nil || prepaySenderDifference[0].Message.Media.ServiceAction.StarGift == nil || @@ -139,10 +143,26 @@ WHERE b.owner_user_id=$1 AND b.box_id=$2`, owner.ID, prepaid.Send.RecipientMessa sharedPrepayMedia.ServiceAction.StarGift == nil || sharedPrepayMedia.ServiceAction.StarGift.GiftMsgID != 0 { t.Fatalf("shared prepay media retained account-local gift_msg_id: media=%+v err=%v", sharedPrepayMedia, err) } - upgraded, err := upgrades.UpgradeStarGift(ctx, domain.StarGiftUpgradeRequest{ - UserID: owner.ID, Ref: domain.SavedStarGiftRef{Owner: ownerPeer, MsgID: purchased.Saved.MsgID}, + if byPrepay, found, err := gifts.GetByRef(ctx, domain.SavedStarGiftRef{Owner: ownerPeer, MsgID: prepaid.Send.RecipientMessage.ID}); err != nil || !found || byPrepay.ID != purchased.Saved.ID { + t.Fatalf("prepay owner message ref = %+v found=%v err=%v", byPrepay, found, err) + } + var ownerPrepayAlias, payerPrepayAlias int + if err := pool.QueryRow(ctx, `SELECT COUNT(*) FROM star_gift_user_message_refs +WHERE owner_user_id=$1 AND msg_id=$2 AND saved_gift_id=$3`, owner.ID, prepaid.Send.RecipientMessage.ID, purchased.Saved.ID).Scan(&ownerPrepayAlias); err != nil { + t.Fatalf("load owner prepay alias: %v", err) + } + if err := pool.QueryRow(ctx, `SELECT COUNT(*) FROM star_gift_user_message_refs +WHERE owner_user_id=$1 AND msg_id=$2`, prepayPayer.ID, prepaid.Send.SenderMessage.ID).Scan(&payerPrepayAlias); err != nil { + t.Fatalf("load payer prepay alias: %v", err) + } + if ownerPrepayAlias != 1 || payerPrepayAlias != 0 { + t.Fatalf("prepay aliases owner=%d payer=%d, want owner-only", ownerPrepayAlias, payerPrepayAlias) + } + upgradeReq := domain.StarGiftUpgradeRequest{ + UserID: owner.ID, Ref: domain.SavedStarGiftRef{Owner: ownerPeer, MsgID: prepaid.Send.RecipientMessage.ID}, RequirePrepaid: true, KeepOriginalDetails: true, CommandKey: "upgrade-" + suffix, Date: now + 2, - }) + } + upgraded, err := upgrades.UpgradeStarGift(ctx, upgradeReq) if err != nil { t.Fatalf("upgrade prepaid gift: %v", err) } @@ -198,7 +218,9 @@ WHERE b.owner_user_id=$1 AND b.box_id=$2`, owner.ID, prepaid.Send.RecipientMessa } upgradeAction := upgraded.Send.RecipientMessage.Media.ServiceAction.StarGiftUnique senderUpgradeAction := upgraded.Send.SenderMessage.Media.ServiceAction.StarGiftUnique - ownerSourceEdit := upgradedSourceEditForUser(upgraded, owner.ID) + ownerSourceEdit := upgradedSourceEditForMessage(upgraded, owner.ID, purchased.Saved.MsgID) + ownerPrepayEdit := upgradedSourceEditForMessage(upgraded, owner.ID, prepaid.Send.RecipientMessage.ID) + payerPrepayEdit := upgradedSourceEditForMessage(upgraded, prepayPayer.ID, prepaid.Send.SenderMessage.ID) if upgradeAction == nil || upgradeAction.SavedID != 0 || upgradeAction.Peer.Type != "" || upgradeAction.Peer.ID != 0 || upgradeAction.CanCraftAt != now+2 || senderUpgradeAction == nil || senderUpgradeAction.SavedID != 0 || senderUpgradeAction.CanCraftAt != now+2 || @@ -208,6 +230,36 @@ WHERE b.owner_user_id=$1 AND b.box_id=$2`, owner.ID, prepaid.Send.RecipientMessa ownerSourceEdit.Message.Media.ServiceAction.StarGift.CanUpgrade { t.Fatalf("upgrade message linkage = action %+v source edit %+v", upgradeAction, ownerSourceEdit) } + if ownerPrepayEdit.Event.Pts <= ownerSourceEdit.Event.Pts || ownerPrepayEdit.Message.Media == nil || + ownerPrepayEdit.Message.Media.ServiceAction == nil || ownerPrepayEdit.Message.Media.ServiceAction.StarGift == nil || + ownerPrepayEdit.Message.Media.ServiceAction.StarGift.CanUpgrade || + ownerPrepayEdit.Message.Media.ServiceAction.StarGift.UpgradeMsgID != upgraded.Send.RecipientMessage.ID { + t.Fatalf("owner prepay card did not converge with upgrade: %+v", ownerPrepayEdit) + } + if payerPrepayEdit.Event.Pts <= prepaid.Send.SenderMessage.Pts || payerPrepayEdit.Message.Media == nil || + payerPrepayEdit.Message.Media.ServiceAction == nil || payerPrepayEdit.Message.Media.ServiceAction.StarGift == nil || + payerPrepayEdit.Message.Media.ServiceAction.StarGift.CanUpgrade || + payerPrepayEdit.Message.Media.ServiceAction.StarGift.UpgradeMsgID != 0 { + t.Fatalf("third-party payer prepay card retained an owner action/link: %+v", payerPrepayEdit) + } + ownerUpgradeDifference, err := NewUpdateEventStore(pool).ListAfter(ctx, owner.ID, upgraded.Send.RecipientMessage.Pts-1, 3) + if err != nil || len(ownerUpgradeDifference) != 3 || + ownerUpgradeDifference[0].Type != domain.UpdateEventNewMessage || + ownerUpgradeDifference[1].Type != domain.UpdateEventEditMessage || ownerUpgradeDifference[1].Message.ID != purchased.Saved.MsgID || + ownerUpgradeDifference[2].Type != domain.UpdateEventEditMessage || ownerUpgradeDifference[2].Message.ID != prepaid.Send.RecipientMessage.ID { + t.Fatalf("owner prepaid upgrade difference = %+v err=%v", ownerUpgradeDifference, err) + } + payerUpgradeDifference, err := NewUpdateEventStore(pool).ListAfter(ctx, prepayPayer.ID, prepaid.Send.SenderMessage.Pts, 1) + if err != nil || len(payerUpgradeDifference) != 1 || payerUpgradeDifference[0].Type != domain.UpdateEventEditMessage || + payerUpgradeDifference[0].Message.ID != prepaid.Send.SenderMessage.ID { + t.Fatalf("payer prepaid upgrade difference = %+v err=%v", payerUpgradeDifference, err) + } + replayedUpgrade, err := upgrades.UpgradeStarGift(ctx, upgradeReq) + if err != nil || !replayedUpgrade.Duplicate || replayedUpgrade.Unique.ID != upgraded.Unique.ID || + upgradedSourceEditForMessage(replayedUpgrade, owner.ID, purchased.Saved.MsgID).Event.Pts != ownerSourceEdit.Event.Pts || + upgradedSourceEditForMessage(replayedUpgrade, owner.ID, prepaid.Send.RecipientMessage.ID).Event.Pts != ownerPrepayEdit.Event.Pts { + t.Fatalf("replay prepaid upgrade from notification = %+v err=%v", replayedUpgrade, err) + } var sharedUpgradeSourceMediaJSON string if err := pool.QueryRow(ctx, `SELECT p.media::text FROM private_messages p JOIN message_boxes b ON b.message_sender_id=p.sender_user_id AND b.private_message_id=p.id @@ -221,6 +273,23 @@ WHERE b.owner_user_id=$1 AND b.box_id=$2`, owner.ID, purchased.Saved.MsgID).Scan sharedUpgradeSourceMedia.ServiceAction.StarGift.GiftMsgID != 0 { t.Fatalf("shared upgraded source media retained account-local message id: media=%+v err=%v", sharedUpgradeSourceMedia, err) } + var sharedUpgradedPrepayMediaJSON string + if err := pool.QueryRow(ctx, `SELECT p.media::text FROM private_messages p +JOIN message_boxes b ON b.message_sender_id=p.sender_user_id AND b.private_message_id=p.id +WHERE b.owner_user_id=$1 AND b.box_id=$2`, owner.ID, prepaid.Send.RecipientMessage.ID).Scan(&sharedUpgradedPrepayMediaJSON); err != nil { + t.Fatalf("load shared upgraded prepay media: %v", err) + } + sharedUpgradedPrepayMedia, err := decodeMessageMedia(sharedUpgradedPrepayMediaJSON) + if err != nil || sharedUpgradedPrepayMedia == nil || sharedUpgradedPrepayMedia.ServiceAction == nil || + sharedUpgradedPrepayMedia.ServiceAction.StarGift == nil || + sharedUpgradedPrepayMedia.ServiceAction.StarGift.CanUpgrade || + sharedUpgradedPrepayMedia.ServiceAction.StarGift.PrepaidUpgradeHash != "" || + sharedUpgradedPrepayMedia.ServiceAction.StarGift.UpgradeMsgID != 0 || + sharedUpgradedPrepayMedia.ServiceAction.StarGift.GiftMsgID != 0 { + t.Fatalf("shared upgraded prepay media retained an action or account-local id: media=%+v err=%v", sharedUpgradedPrepayMedia, err) + } + verifyPrepaidMessageRefMigration(t, ctx, pool, purchased.Saved.ID, owner.ID, prepayPayer.ID, + prepaid.Send.RecipientMessage.ID, prepaid.Send.SenderMessage.ID, upgraded.Send.RecipientMessage.ID) var sharedUpgradeMediaJSON string if err := pool.QueryRow(ctx, `SELECT p.media::text FROM private_messages p JOIN message_boxes b ON b.message_sender_id=p.sender_user_id AND b.private_message_id=p.id @@ -361,6 +430,16 @@ WHERE b.owner_user_id=$1 AND b.box_id=$2`, owner.ID, upgraded.Send.RecipientMess if err != nil || transferred.Unique.Owner != ownerPeer || transferred.Saved.TransferStars != 25 || transferred.Balance.Balance != 9975 { t.Fatalf("paid transfer = %+v err %v", transferred, err) } + const historicalOwnerMessageID = 2_147_483_000 + if _, err := pool.Exec(ctx, `INSERT INTO star_gift_user_message_refs(owner_user_id,msg_id,saved_gift_id) +VALUES($1,$2,$3)`, resaleBuyer.ID, historicalOwnerMessageID, transferred.Saved.ID); err != nil { + t.Fatalf("insert historical old-owner message ref: %v", err) + } + if _, err := gifts.ResolveSavedIDs(ctx, ownerPeer, []domain.SavedStarGiftRef{{ + Owner: ownerPeer, MsgID: historicalOwnerMessageID, + }}); !errors.Is(err, domain.ErrStarGiftNotFound) { + t.Fatalf("current owner resolved another owner's historical message ref: %v", err) + } var retiredSourceMediaJSON string var retiredSourcePTS int if err := pool.QueryRow(ctx, `SELECT media::text,pts FROM message_boxes @@ -1215,6 +1294,100 @@ func issueLifecyclePurchaseForm(t *testing.T, ctx context.Context, lifecycle *St return req } +func upgradedSourceEditForMessage(result domain.StarGiftUpgradeResult, userID int64, messageID int) domain.EditedMessageForUser { + for _, edit := range result.SourceEdits { + if edit.UserID == userID && edit.Message.ID == messageID { + return edit + } + } + return domain.EditedMessageForUser{UserID: userID} +} + +func verifyPrepaidMessageRefMigration( + t *testing.T, + ctx context.Context, + pool *pgxpool.Pool, + savedGiftID int64, + ownerUserID int64, + payerUserID int64, + ownerPrepayMessageID int, + payerPrepayMessageID int, + ownerUpgradeMessageID int, +) { + t.Helper() + tx, err := pool.Begin(ctx) + if err != nil { + t.Fatalf("begin prepaid message migration probe: %v", err) + } + defer func() { _ = tx.Rollback(context.Background()) }() + + var messageSenderID, privateMessageID int64 + if err := tx.QueryRow(ctx, `SELECT message_sender_id,private_message_id FROM message_boxes +WHERE owner_user_id=$1 AND box_id=$2 AND NOT deleted`, ownerUserID, ownerPrepayMessageID). + Scan(&messageSenderID, &privateMessageID); err != nil { + t.Fatalf("load prepaid message root for migration probe: %v", err) + } + if _, err := tx.Exec(ctx, `DELETE FROM star_gift_user_message_refs +WHERE owner_user_id=$1 AND msg_id=$2 AND saved_gift_id=$3`, ownerUserID, ownerPrepayMessageID, savedGiftID); err != nil { + t.Fatalf("remove prepaid alias for migration probe: %v", err) + } + if _, err := tx.Exec(ctx, `UPDATE message_boxes +SET media=jsonb_set(media #- '{service_action,star_gift,upgrade_msg_id}', + '{service_action,star_gift,can_upgrade}','true'::jsonb,true) +WHERE message_sender_id=$1 AND private_message_id=$2 AND NOT deleted`, messageSenderID, privateMessageID); err != nil { + t.Fatalf("restore stale prepaid message boxes for migration probe: %v", err) + } + if _, err := tx.Exec(ctx, `UPDATE private_messages +SET media=jsonb_set(media #- '{service_action,star_gift,upgrade_msg_id}', + '{service_action,star_gift,can_upgrade}','true'::jsonb,true) +WHERE sender_user_id=$1 AND id=$2`, messageSenderID, privateMessageID); err != nil { + t.Fatalf("restore stale shared prepaid message for migration probe: %v", err) + } + + migrationSQL, err := deploy.Migrations.ReadFile("migrations/0135_star_gift_prepaid_message_refs.up.sql") + if err != nil { + t.Fatalf("read prepaid message migration: %v", err) + } + if _, err := tx.Exec(ctx, string(migrationSQL)); err != nil { + t.Fatalf("apply prepaid message migration probe: %v", err) + } + + var aliasCount int + if err := tx.QueryRow(ctx, `SELECT COUNT(*) FROM star_gift_user_message_refs +WHERE owner_user_id=$1 AND msg_id=$2 AND saved_gift_id=$3`, ownerUserID, ownerPrepayMessageID, savedGiftID).Scan(&aliasCount); err != nil || aliasCount != 1 { + t.Fatalf("migrated prepaid alias count=%d err=%v", aliasCount, err) + } + assertMigratedAction := func(userID int64, messageID int, wantUpgradeMessageID int) { + t.Helper() + var mediaJSON string + var pts int + if err := tx.QueryRow(ctx, `SELECT media::text,pts FROM message_boxes +WHERE owner_user_id=$1 AND box_id=$2 AND NOT deleted`, userID, messageID).Scan(&mediaJSON, &pts); err != nil { + t.Fatalf("load migrated prepaid box %d/%d: %v", userID, messageID, err) + } + media, err := decodeMessageMedia(mediaJSON) + if err != nil || media == nil || media.ServiceAction == nil || media.ServiceAction.StarGift == nil || + media.ServiceAction.StarGift.CanUpgrade || media.ServiceAction.StarGift.PrepaidUpgradeHash != "" || + media.ServiceAction.StarGift.UpgradeMsgID != wantUpgradeMessageID { + t.Fatalf("migrated prepaid box %d/%d = %+v err=%v", userID, messageID, media, err) + } + var eventCount, outboxCount int + if err := tx.QueryRow(ctx, `SELECT COUNT(*) FROM user_update_events +WHERE user_id=$1 AND pts=$2 AND event_type='edit_message' AND message_box_id=$3`, userID, pts, messageID).Scan(&eventCount); err != nil { + t.Fatalf("load migrated prepaid event: %v", err) + } + if err := tx.QueryRow(ctx, `SELECT COUNT(*) FROM dispatch_outbox +WHERE target_user_id=$1 AND pts=$2 AND event_type='edit_message'`, userID, pts).Scan(&outboxCount); err != nil { + t.Fatalf("load migrated prepaid outbox: %v", err) + } + if eventCount != 1 || outboxCount != 1 { + t.Fatalf("migrated prepaid event/outbox counts=%d/%d", eventCount, outboxCount) + } + } + assertMigratedAction(ownerUserID, ownerPrepayMessageID, ownerUpgradeMessageID) + assertMigratedAction(payerUserID, payerPrepayMessageID, 0) +} + func craftedSourceEditForUserAndGift(result domain.StarGiftCraftResult, userID, uniqueGiftID int64) domain.EditedMessageForUser { for _, edit := range result.SourceEdits { if edit.UserID != userID { diff --git a/internal/store/postgres/star_gift_lifecycle_migration_integration_test.go b/internal/store/postgres/star_gift_lifecycle_migration_integration_test.go index 5000638f..ebcfe739 100644 --- a/internal/store/postgres/star_gift_lifecycle_migration_integration_test.go +++ b/internal/store/postgres/star_gift_lifecycle_migration_integration_test.go @@ -14,7 +14,7 @@ func TestStarGiftLifecycleMigrationsApply(t *testing.T) { if err != nil { t.Fatalf("migrate star gift lifecycle schema: %v", err) } - if status.Dirty || status.Empty || status.Version != 134 { - t.Fatalf("migration status = %+v, want clean version 134", status) + if status.Dirty || status.Empty || status.Version != 135 { + t.Fatalf("migration status = %+v, want clean version 135", status) } } diff --git a/internal/store/postgres/star_gift_upgrade.go b/internal/store/postgres/star_gift_upgrade.go index b786e289..31ea25e9 100644 --- a/internal/store/postgres/star_gift_upgrade.go +++ b/internal/store/postgres/star_gift_upgrade.go @@ -341,11 +341,45 @@ func starGiftUpgradeUniqueAction(saved domain.SavedStarGift, unique domain.Uniqu } } -// markPrivateStarGiftSourceUpgradedTx rewrites both visible copies of the -// original gift service message in the same transaction that creates the -// unique gift message. upgrade_msg_id is box-local, so each owner projection -// must point at that owner's copy of the new service message. Every rewrite is -// a durable edit_message event with its own pts and outbox row. +func userStarGiftSourceMessageIDs(ctx context.Context, db interface { + Query(context.Context, string, ...any) (pgx.Rows, error) +}, saved domain.SavedStarGift) ([]int, error) { + if saved.Owner.Type != domain.PeerTypeUser || saved.Owner.ID <= 0 || saved.ID <= 0 || saved.MsgID <= 0 { + return nil, domain.ErrStarGiftCollectibleInvalid + } + messageIDs := []int{saved.MsgID} + rows, err := db.Query(ctx, ` +SELECT msg_id FROM star_gift_user_message_refs +WHERE owner_user_id=$1 AND saved_gift_id=$2 AND msg_id<>$3 +ORDER BY msg_id`, saved.Owner.ID, saved.ID, saved.MsgID) + if err != nil { + return nil, fmt.Errorf("list star gift source message refs: %w", err) + } + defer rows.Close() + for rows.Next() { + var msgID int + if err := rows.Scan(&msgID); err != nil { + return nil, fmt.Errorf("scan star gift source message ref: %w", err) + } + if msgID <= 0 { + return nil, fmt.Errorf("star gift source message ref has invalid id") + } + messageIDs = append(messageIDs, msgID) + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("iterate star gift source message refs: %w", err) + } + return messageIDs, nil +} + +// markPrivateStarGiftSourceUpgradedTx rewrites every ordinary gift projection +// owned by the source aggregate: the original gift message and each separately +// prepaid-upgrade notification. The two visible boxes of every logical private +// message are updated together. upgrade_msg_id is box-local and is set only +// when that viewer owns a box for the emitted unique-gift message; a third-party +// payer sees the prepayment become non-actionable without receiving an invalid +// owner-local link. Every rewrite is a durable edit_message event with its own +// pts and outbox row. func (s *StarGiftUpgradeStore) markPrivateStarGiftSourceUpgradedTx( ctx context.Context, tx pgx.Tx, @@ -356,30 +390,11 @@ func (s *StarGiftUpgradeStore) markPrivateStarGiftSourceUpgradedTx( if saved.Owner.Type != domain.PeerTypeUser || saved.Owner.ID != req.UserID || saved.MsgID <= 0 { return nil, domain.ErrStarGiftCollectibleInvalid } + messageIDs, err := userStarGiftSourceMessageIDs(ctx, tx, saved) + if err != nil { + return nil, err + } q := sqlcgen.New(tx) - target, err := q.GetMessageBoxForEdit(ctx, sqlcgen.GetMessageBoxForEditParams{ - OwnerUserID: req.UserID, - BoxID: int32(saved.MsgID), - PeerType: string(domain.PeerTypeUser), - PeerID: saved.FromUserID, - }) - if err != nil { - if errors.Is(err, pgx.ErrNoRows) { - return nil, domain.ErrStarGiftCollectibleInvalid - } - return nil, fmt.Errorf("lock star gift source message: %w", err) - } - boxes, err := q.ListVisibleMessageBoxesByPrivateMessage(ctx, sqlcgen.ListVisibleMessageBoxesByPrivateMessageParams{ - OwnerUserIds: privateMessageOwnerIDs(req.UserID, saved.FromUserID), - MessageSenderID: target.MessageSenderID, - PrivateMessageID: target.PrivateMessageID, - }) - if err != nil { - return nil, fmt.Errorf("list star gift source message boxes: %w", err) - } - if len(boxes) == 0 { - return nil, domain.ErrStarGiftCollectibleInvalid - } upgradeMessageIDs := make(map[int64]int, 2) if sent.SenderMessage.OwnerUserID > 0 && sent.SenderMessage.ID > 0 { upgradeMessageIDs[sent.SenderMessage.OwnerUserID] = sent.SenderMessage.ID @@ -387,87 +402,158 @@ func (s *StarGiftUpgradeStore) markPrivateStarGiftSourceUpgradedTx( if sent.RecipientMessage.OwnerUserID > 0 && sent.RecipientMessage.ID > 0 { upgradeMessageIDs[sent.RecipientMessage.OwnerUserID] = sent.RecipientMessage.ID } - edits := make([]domain.EditedMessageForUser, 0, len(boxes)) - var privateMediaJSON []byte - for _, box := range boxes { - upgradeMessageID := upgradeMessageIDs[box.OwnerUserID] - if upgradeMessageID <= 0 { - return nil, fmt.Errorf("upgrade service message missing box for user %d", box.OwnerUserID) + edits := make([]domain.EditedMessageForUser, 0, len(messageIDs)*2) + seenPrivateMessages := make(map[string]struct{}, len(messageIDs)) + primaryRewritten := false + for _, sourceMessageID := range messageIDs { + var peerType string + var peerID int64 + err := tx.QueryRow(ctx, ` +SELECT peer_type,peer_id FROM message_boxes +WHERE owner_user_id=$1 AND box_id=$2 AND NOT deleted +FOR UPDATE`, req.UserID, sourceMessageID).Scan(&peerType, &peerID) + if errors.Is(err, pgx.ErrNoRows) { + if sourceMessageID == saved.MsgID { + return nil, domain.ErrStarGiftCollectibleInvalid + } + continue } - media, err := decodeMessageMedia(box.MediaJson) if err != nil { - return nil, fmt.Errorf("decode star gift source media: %w", err) + return nil, fmt.Errorf("lock star gift source ref %d: %w", sourceMessageID, err) } - if media == nil || media.Kind != domain.MessageMediaKindService || media.ServiceAction == nil || - media.ServiceAction.Kind != domain.MessageServiceActionStarGift || media.ServiceAction.StarGift == nil { - return nil, fmt.Errorf("star gift source message %d has invalid media", box.BoxID) + if peerType != string(domain.PeerTypeUser) || peerID <= 0 { + return nil, fmt.Errorf("star gift source ref %d is not private", sourceMessageID) } - action := media.ServiceAction.StarGift - if action.UpgradeMsgID != 0 && action.UpgradeMsgID != upgradeMessageID { - return nil, fmt.Errorf("star gift source message %d has conflicting upgrade message %d", box.BoxID, action.UpgradeMsgID) - } - action.UpgradeMsgID = upgradeMessageID - action.CanUpgrade = false - mediaJSON, err := encodeMessageMedia(media) + target, err := q.GetMessageBoxForEdit(ctx, sqlcgen.GetMessageBoxForEditParams{ + OwnerUserID: req.UserID, BoxID: int32(sourceMessageID), PeerType: peerType, PeerID: peerID, + }) if err != nil { - return nil, fmt.Errorf("encode upgraded star gift source media: %w", err) + return nil, fmt.Errorf("load star gift source ref %d: %w", sourceMessageID, err) } - pts, err := s.messages.reservePts(ctx, tx, box.OwnerUserID) + ownerMedia, err := decodeMessageMedia(target.MediaJson) if err != nil { - return nil, fmt.Errorf("allocate star gift source edit pts: %w", err) + return nil, fmt.Errorf("decode star gift source ref %d: %w", sourceMessageID, err) } - tag, err := tx.Exec(ctx, ` + ownerAction := privateStarGiftAction(ownerMedia) + if ownerAction == nil { + // The newly emitted unique action is registered before source edits in + // the same transaction and belongs to the same aggregate, but is not a + // source projection to rewrite. + if privateStarGiftUniqueAction(ownerMedia) != nil { + continue + } + return nil, fmt.Errorf("star gift source ref %d has invalid media", sourceMessageID) + } + if ownerAction.GiftID != saved.GiftID { + return nil, fmt.Errorf("star gift source ref %d points to gift %d", sourceMessageID, ownerAction.GiftID) + } + if sourceMessageID != saved.MsgID && (!ownerAction.UpgradeSeparate || !ownerAction.PrepaidUpgrade || ownerAction.GiftMsgID != saved.MsgID) { + return nil, fmt.Errorf("star gift source ref %d is not a prepaid notification for message %d", sourceMessageID, saved.MsgID) + } + logicalKey := fmt.Sprintf("%d:%d", target.MessageSenderID, target.PrivateMessageID) + if _, duplicate := seenPrivateMessages[logicalKey]; duplicate { + continue + } + seenPrivateMessages[logicalKey] = struct{}{} + boxes, err := q.ListVisibleMessageBoxesByPrivateMessage(ctx, sqlcgen.ListVisibleMessageBoxesByPrivateMessageParams{ + OwnerUserIds: privateMessageOwnerIDs(req.UserID, peerID), MessageSenderID: target.MessageSenderID, + PrivateMessageID: target.PrivateMessageID, + }) + if err != nil { + return nil, fmt.Errorf("list star gift source ref %d boxes: %w", sourceMessageID, err) + } + if len(boxes) == 0 { + return nil, domain.ErrStarGiftCollectibleInvalid + } + var privateMediaJSON []byte + for _, box := range boxes { + media, err := decodeMessageMedia(box.MediaJson) + if err != nil { + return nil, fmt.Errorf("decode star gift source media: %w", err) + } + action := privateStarGiftAction(media) + if action == nil || action.GiftID != saved.GiftID { + return nil, fmt.Errorf("star gift source message %d has invalid media", box.BoxID) + } + upgradeMessageID := upgradeMessageIDs[box.OwnerUserID] + if action.UpgradeMsgID != 0 && upgradeMessageID > 0 && action.UpgradeMsgID != upgradeMessageID { + return nil, fmt.Errorf("star gift source message %d has conflicting upgrade message %d", box.BoxID, action.UpgradeMsgID) + } + if upgradeMessageID > 0 { + action.UpgradeMsgID = upgradeMessageID + } else { + if box.OwnerUserID == req.UserID { + return nil, fmt.Errorf("upgrade service message missing owner box") + } + action.UpgradeMsgID = 0 + } + action.CanUpgrade = false + action.PrepaidUpgradeHash = "" + mediaJSON, err := encodeMessageMedia(media) + if err != nil { + return nil, fmt.Errorf("encode upgraded star gift source media: %w", err) + } + pts, err := s.messages.reservePts(ctx, tx, box.OwnerUserID) + if err != nil { + return nil, fmt.Errorf("allocate star gift source edit pts: %w", err) + } + tag, err := tx.Exec(ctx, ` UPDATE message_boxes SET media=$3, pts=$4 WHERE owner_user_id=$1 AND box_id=$2 AND NOT deleted`, box.OwnerUserID, box.BoxID, mediaJSON, int32(pts)) - if err != nil { - return nil, fmt.Errorf("update star gift source message box: %w", err) - } - if tag.RowsAffected() != 1 { - return nil, fmt.Errorf("update star gift source message box lost row") - } - msg, err := messageFromVisibleBoxRow(box) - if err != nil { - return nil, err - } - msg.Media = media - msg.Pts = pts - if err := replaceMessageBoxMediaIndexTx(ctx, tx, msg.OwnerUserID, msg.Peer.ID, msg.ID, msg.Date, msg.Media, msg.Entities); err != nil { - return nil, err - } - event := domain.UpdateEvent{ - UserID: msg.OwnerUserID, Type: domain.UpdateEventEditMessage, - Pts: pts, PtsCount: 1, Date: req.Date, Message: msg, - } - if err := appendUserUpdateEvent(ctx, tx, q, msg.OwnerUserID, event); err != nil { - return nil, fmt.Errorf("append star gift source edit event: %w", err) - } - dispatchAuthKeyID := [8]byte{} - dispatchSessionID := int64(0) - if msg.OwnerUserID == req.UserID { - dispatchAuthKeyID = req.OriginAuthKeyID - dispatchSessionID = req.OriginSessionID - } - if err := enqueueDispatch(ctx, q, sqlcgen.EnqueueDispatchParams{ - TargetUserID: msg.OwnerUserID, Pts: int32(pts), EventType: string(domain.UpdateEventEditMessage), - ExcludeAuthKeyID: authKeyIDToInt64(dispatchAuthKeyID), ExcludeSessionID: dispatchSessionID, - }); err != nil { - return nil, fmt.Errorf("enqueue star gift source edit: %w", err) - } - if box.OwnerUserID == box.MessageSenderID || len(privateMediaJSON) == 0 { - privateMediaJSON, err = encodeSharedPrivateStarGiftMedia(media) + if err != nil { + return nil, fmt.Errorf("update star gift source message box: %w", err) + } + if tag.RowsAffected() != 1 { + return nil, fmt.Errorf("update star gift source message box lost row") + } + msg, err := messageFromVisibleBoxRow(box) if err != nil { return nil, err } + msg.Media = media + msg.Pts = pts + if err := replaceMessageBoxMediaIndexTx(ctx, tx, msg.OwnerUserID, msg.Peer.ID, msg.ID, msg.Date, msg.Media, msg.Entities); err != nil { + return nil, err + } + event := domain.UpdateEvent{UserID: msg.OwnerUserID, Type: domain.UpdateEventEditMessage, + Pts: pts, PtsCount: 1, Date: req.Date, Message: msg} + if err := appendUserUpdateEvent(ctx, tx, q, msg.OwnerUserID, event); err != nil { + return nil, fmt.Errorf("append star gift source edit event: %w", err) + } + dispatchAuthKeyID := [8]byte{} + dispatchSessionID := int64(0) + if msg.OwnerUserID == req.UserID { + dispatchAuthKeyID = req.OriginAuthKeyID + dispatchSessionID = req.OriginSessionID + } + if err := enqueueDispatch(ctx, q, sqlcgen.EnqueueDispatchParams{ + TargetUserID: msg.OwnerUserID, Pts: int32(pts), EventType: string(domain.UpdateEventEditMessage), + ExcludeAuthKeyID: authKeyIDToInt64(dispatchAuthKeyID), ExcludeSessionID: dispatchSessionID, + }); err != nil { + return nil, fmt.Errorf("enqueue star gift source edit: %w", err) + } + if box.OwnerUserID == box.MessageSenderID || len(privateMediaJSON) == 0 { + privateMediaJSON, err = encodeSharedPrivateStarGiftMedia(media) + if err != nil { + return nil, err + } + } + edits = append(edits, domain.EditedMessageForUser{UserID: msg.OwnerUserID, Message: msg, Event: event}) } - edits = append(edits, domain.EditedMessageForUser{UserID: msg.OwnerUserID, Message: msg, Event: event}) - } - if len(privateMediaJSON) == 0 { - return nil, fmt.Errorf("upgrade source message missing private media projection") - } - if _, err := tx.Exec(ctx, ` + if len(privateMediaJSON) == 0 { + return nil, fmt.Errorf("upgrade source message missing private media projection") + } + if _, err := tx.Exec(ctx, ` UPDATE private_messages SET media=$3 WHERE sender_user_id=$1 AND id=$2`, target.MessageSenderID, target.PrivateMessageID, privateMediaJSON); err != nil { - return nil, fmt.Errorf("update star gift source private message: %w", err) + return nil, fmt.Errorf("update star gift source private message: %w", err) + } + if sourceMessageID == saved.MsgID { + primaryRewritten = true + } + } + if !primaryRewritten { + return nil, domain.ErrStarGiftCollectibleInvalid } return edits, nil } @@ -678,46 +764,77 @@ func (s *StarGiftUpgradeStore) loadUpgradeSourceReplay(ctx context.Context, req if pts <= 0 || saved.MsgID <= 0 { return nil, domain.ErrStarGiftCollectibleInvalid } - var privateMessageID, messageSenderID int64 - err := s.db.QueryRow(ctx, ` -SELECT private_message_id,message_sender_id FROM message_boxes -WHERE owner_user_id=$1 AND box_id=$2 AND peer_type='user' AND peer_id=$3 AND NOT deleted`, - req.UserID, saved.MsgID, saved.FromUserID).Scan(&privateMessageID, &messageSenderID) - if errors.Is(err, pgx.ErrNoRows) { - // A later delete event is authoritative; replaying the old edit here - // would transiently resurrect the source message. - return nil, nil - } - if err != nil { - return nil, fmt.Errorf("load star gift source replay message: %w", err) - } - boxes, err := sqlcgen.New(s.db).ListVisibleMessageBoxesByPrivateMessage(ctx, sqlcgen.ListVisibleMessageBoxesByPrivateMessageParams{ - OwnerUserIds: []int64{req.UserID}, MessageSenderID: messageSenderID, PrivateMessageID: privateMessageID, - }) - if err != nil { - return nil, fmt.Errorf("load star gift source replay box: %w", err) - } - if len(boxes) != 1 || int(boxes[0].BoxID) != saved.MsgID { - return nil, domain.ErrStarGiftCollectibleInvalid - } - var eventDate int - err = s.db.QueryRow(ctx, ` -SELECT date FROM user_update_events -WHERE user_id=$1 AND pts=$2 AND event_type='edit_message' AND message_box_id=$3`, - req.UserID, pts, saved.MsgID).Scan(&eventDate) - if err != nil { - if errors.Is(err, pgx.ErrNoRows) { - return nil, domain.ErrStarGiftCollectibleInvalid - } - return nil, fmt.Errorf("load star gift source replay event: %w", err) - } - msg, err := messageFromVisibleBoxRow(boxes[0]) + messageIDs, err := userStarGiftSourceMessageIDs(ctx, s.db, saved) if err != nil { return nil, err } - msg.Pts = pts - event := domain.UpdateEvent{UserID: req.UserID, Type: domain.UpdateEventEditMessage, Pts: pts, PtsCount: 1, Date: eventDate, Message: msg} - return []domain.EditedMessageForUser{{UserID: req.UserID, Message: msg, Event: event}}, nil + edits := make([]domain.EditedMessageForUser, 0, len(messageIDs)) + for _, messageID := range messageIDs { + var privateMessageID, messageSenderID int64 + err := s.db.QueryRow(ctx, ` +SELECT private_message_id,message_sender_id FROM message_boxes +WHERE owner_user_id=$1 AND box_id=$2 AND peer_type='user' AND NOT deleted`, + req.UserID, messageID).Scan(&privateMessageID, &messageSenderID) + if errors.Is(err, pgx.ErrNoRows) { + // A later delete event is authoritative; replaying the old edit here + // would transiently resurrect that source projection. + continue + } + if err != nil { + return nil, fmt.Errorf("load star gift source replay message %d: %w", messageID, err) + } + boxes, err := sqlcgen.New(s.db).ListVisibleMessageBoxesByPrivateMessage(ctx, sqlcgen.ListVisibleMessageBoxesByPrivateMessageParams{ + OwnerUserIds: []int64{req.UserID}, MessageSenderID: messageSenderID, PrivateMessageID: privateMessageID, + }) + if err != nil { + return nil, fmt.Errorf("load star gift source replay box %d: %w", messageID, err) + } + if len(boxes) != 1 || int(boxes[0].BoxID) != messageID { + return nil, domain.ErrStarGiftCollectibleInvalid + } + media, err := decodeMessageMedia(boxes[0].MediaJson) + if err != nil { + return nil, fmt.Errorf("decode star gift source replay box %d: %w", messageID, err) + } + action := privateStarGiftAction(media) + if action == nil { + if privateStarGiftUniqueAction(media) != nil { + continue + } + return nil, fmt.Errorf("star gift source replay box %d has invalid media", messageID) + } + if action.GiftID != saved.GiftID || action.CanUpgrade || action.UpgradeMsgID != saved.UpgradeMsgID { + return nil, domain.ErrStarGiftCollectibleInvalid + } + var eventPts, eventDate int + if messageID == saved.MsgID { + eventPts = pts + err = s.db.QueryRow(ctx, ` +SELECT date FROM user_update_events +WHERE user_id=$1 AND pts=$2 AND event_type='edit_message' AND message_box_id=$3`, + req.UserID, eventPts, messageID).Scan(&eventDate) + } else { + err = s.db.QueryRow(ctx, ` +SELECT pts,date FROM user_update_events +WHERE user_id=$1 AND pts>=$2 AND event_type='edit_message' AND message_box_id=$3 +ORDER BY pts LIMIT 1`, req.UserID, pts, messageID).Scan(&eventPts, &eventDate) + } + if err != nil { + if errors.Is(err, pgx.ErrNoRows) { + return nil, domain.ErrStarGiftCollectibleInvalid + } + return nil, fmt.Errorf("load star gift source replay event %d: %w", messageID, err) + } + msg, err := messageFromVisibleBoxRow(boxes[0]) + if err != nil { + return nil, err + } + msg.Pts = eventPts + event := domain.UpdateEvent{UserID: req.UserID, Type: domain.UpdateEventEditMessage, + Pts: eventPts, PtsCount: 1, Date: eventDate, Message: msg} + edits = append(edits, domain.EditedMessageForUser{UserID: req.UserID, Message: msg, Event: event}) + } + return edits, nil } func (s *StarGiftUpgradeStore) StarGiftUpgradeReceipt(ctx context.Context, userID int64, commandKey string) (domain.StarGiftUpgradeReceipt, bool, error) { diff --git a/internal/store/postgres/star_gift_user_message_ref.go b/internal/store/postgres/star_gift_user_message_ref.go index b1a4da81..e69898ec 100644 --- a/internal/store/postgres/star_gift_user_message_ref.go +++ b/internal/store/postgres/star_gift_user_message_ref.go @@ -9,9 +9,11 @@ import ( // registerUserStarGiftMessageRef records an owner-scoped service-message alias // for a user-owned gift. Official clients may continue from a freshly emitted -// messageActionStarGiftUnique and pass that message id to a lifecycle RPC, -// while payments.getSavedStarGifts may still expose the original received gift -// message as the aggregate's primary msg_id. +// messageActionStarGiftUnique or a separate prepaid-upgrade notification and +// pass that message id to a lifecycle RPC, while payments.getSavedStarGifts may +// still expose the original received gift message as the aggregate's primary +// msg_id. expectedUniqueGiftID is zero for an ordinary gift and positive for a +// unique gift; the write boundary never aliases across lifecycle states. func registerUserStarGiftMessageRef( ctx context.Context, tx pgx.Tx, @@ -20,7 +22,7 @@ func registerUserStarGiftMessageRef( savedGiftID int64, uniqueGiftID int64, ) error { - if ownerUserID <= 0 || msgID <= 0 || savedGiftID <= 0 || uniqueGiftID <= 0 { + if ownerUserID <= 0 || msgID <= 0 || savedGiftID <= 0 || uniqueGiftID < 0 { return fmt.Errorf("register user star gift message ref: invalid identity") } tag, err := tx.Exec(ctx, ` @@ -28,7 +30,8 @@ INSERT INTO star_gift_user_message_refs(owner_user_id,msg_id,saved_gift_id) SELECT $1,$2,p.id FROM peer_star_gifts p WHERE p.id=$3 AND p.owner_peer_type='user' AND p.owner_peer_id=$1 - AND p.unique_gift_id=$4 AND p.lifecycle_status='active' + AND (($4::bigint=0 AND p.unique_gift_id IS NULL) OR ($4::bigint>0 AND p.unique_gift_id=$4::bigint)) + AND p.lifecycle_status='active' ON CONFLICT(owner_user_id,msg_id) DO UPDATE SET saved_gift_id=EXCLUDED.saved_gift_id WHERE star_gift_user_message_refs.saved_gift_id=EXCLUDED.saved_gift_id`, ownerUserID, msgID, savedGiftID, uniqueGiftID)