-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathservice.go
More file actions
839 lines (792 loc) · 32 KB
/
Copy pathservice.go
File metadata and controls
839 lines (792 loc) · 32 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
// Package chat implements chatapi.Service — one thread per pair of accounts, its
// messages, and the read marks behind an unread badge.
package chat
import (
"context"
"fmt"
"log/slog"
"strings"
"time"
"github.com/go-playground/validator/v10"
accountapi "shopnexus/internal/module/account/api"
chatapi "shopnexus/internal/module/chat/api"
"shopnexus/internal/module/chat/domain"
"shopnexus/internal/module/chat/port"
"shopnexus/internal/module/common"
"shopnexus/internal/shared/id"
"shopnexus/internal/shared/realtime"
"shopnexus/internal/shared/validation"
)
type Service struct {
repo port.Repository
// accounts answers what a counterparty is called, which is what an inbox row shows.
accounts accountapi.Service
// uploads is this module's own resource table plus the object store. A message
// attachment belongs to the module that took the upload, and resolving one through here
// is what puts a live link on it rather than an id nothing can render.
uploads common.Uploads
v *validator.Validate
log *slog.Logger
// fanout pushes the realtime facts in event.go to the socket a participant may have
// open. Best-effort: a write always commits whether or not anybody is listening.
fanout realtime.Fanout
}
func NewService(repo port.Repository, accounts accountapi.Service, uploads common.Uploads, v *validator.Validate, log *slog.Logger, fanout realtime.Fanout) *Service {
return &Service{repo: repo, accounts: accounts, uploads: uploads, v: v, log: log, fanout: fanout}
}
// notify pushes a fact to one account, best-effort.
//
// A realtime failure never fails the command: the row is committed by the time this runs,
// so the alternative is answering 500 for a write that happened. The client re-reads on
// reconnect, which is what covers a dropped event.
func notify[T any](ctx context.Context, s *Service, accountID int64, e realtime.Event[T], data T) {
if err := realtime.Notify(ctx, s.fanout, accountID, e, data); err != nil {
s.log.Warn("realtime notify failed", "code", e.Code, "account_id", accountID, "err", err)
}
}
var _ chatapi.Service = (*Service)(nil)
// CreateUpload reserves a row and a signed slot for a message attachment. The client PUTs
// the bytes at the store and confirms; until then the resource resolves to nothing, so a
// half-finished upload cannot be attached to a message.
func (s *Service) CreateUpload(ctx context.Context, req common.CreateUploadRequest) (common.UploadSlotDTO, error) {
if err := s.v.Struct(req); err != nil {
return common.UploadSlotDTO{}, err
}
slot, err := s.uploads.Presign(ctx, req.ActorID.Int64(), "message", common.UploadRequest{
Filename: req.Filename, Mime: req.Mime, Size: req.Size,
})
if err != nil {
return common.UploadSlotDTO{}, err
}
return slot.ToDTO(), nil
}
// ConfirmUpload makes the attachment real, with the size the store reports rather than the
// one the client declared. Scoped to the uploader: a resource id is guessable, and
// confirming somebody else's slot would be claiming their upload.
func (s *Service) ConfirmUpload(ctx context.Context, req common.ConfirmUploadRequest) (common.ResourceDTO, error) {
if err := s.v.Struct(req); err != nil {
return common.ResourceDTO{}, err
}
res, err := s.uploads.Confirm(ctx, req.ActorID.Int64(), req.ID.Int64())
if err != nil {
return common.ResourceDTO{}, err
}
return res.ToDTO(), nil
}
// ListConversations is the inbox, latest activity first. Two extra queries for the whole
// page — the last message and the unread counts — rather than per row.
func (s *Service) ListConversations(ctx context.Context, req chatapi.ListConversationsRequest) (chatapi.ConversationPage, error) {
before, beforeID, err := parseCursor(req.Cursor)
if err != nil {
return chatapi.ConversationPage{}, err
}
// One more than asked, so "is there another page" is answered without a count.
threads, err := s.repo.ListConversations(ctx, port.InboxFilter{
AccountID: req.ActorID.Int64(),
Before: before,
BeforeID: beforeID,
Limit: req.Limit + 1,
})
if err != nil {
return chatapi.ConversationPage{}, fmt.Errorf("list conversations: %w", err)
}
hasMore := len(threads) > req.Limit
if hasMore {
threads = threads[:req.Limit]
}
out, err := s.inboxRows(ctx, req.ActorID.Int64(), threads)
if err != nil {
return chatapi.ConversationPage{}, err
}
page := chatapi.ConversationPage{Data: out, Meta: chatapi.CursorInfo{HasMore: hasMore}}
if hasMore && len(threads) > 0 {
last := threads[len(threads)-1]
page.Meta.NextCursor = formatCursor(last.LastMessageAt, last.ID)
}
return page, nil
}
// StartConversation opens the thread with one account, or answers the one that already
// exists — there is one per pair, so the route is idempotent by construction.
func (s *Service) StartConversation(ctx context.Context, req chatapi.StartConversationRequest) (chatapi.Conversation, error) {
// Support is reached by raising a ticket, never by messaging the desk: its id is public — it is
// the counterparty of every ticket thread — and a direct thread with it is one no moderator can
// read, so a user who believed they were writing to support would be writing to nobody.
if desk, err := s.desk(ctx); err == nil && req.AccountID.Int64() == desk {
return chatapi.Conversation{}, domain.ErrConversationWithSupport
}
thread, err := s.repo.EnsureConversation(ctx, req.ActorID.Int64(), req.AccountID.Int64())
if err != nil {
return chatapi.Conversation{}, fmt.Errorf("ensure conversation: %w", err)
}
return s.inboxRow(ctx, req.ActorID.Int64(), thread, nil, 0)
}
// OpenTicketThread is the thread behind a ticket: the requester on one side, the support desk's own
// account on the other. The desk rather than a moderator, so whoever answers stays anonymous and the
// next one inherits the thread instead of starting another.
//
// Idempotent, because it is the second half of a two-schema write: trust's ticket row lands first,
// and a retry — or a repair on the next read — has to find the thread this call already made rather
// than a second one.
func (s *Service) OpenTicketThread(ctx context.Context, req chatapi.OpenTicketThreadRequest) (chatapi.Conversation, error) {
if err := validation.AsError(s.v.Struct(req)); err != nil {
return chatapi.Conversation{}, err
}
desk, err := s.desk(ctx)
if err != nil {
return chatapi.Conversation{}, err
}
thread, err := s.repo.EnsureTicketThread(ctx, req.RequesterID.Int64(), desk, req.TicketID.Int64())
if err != nil {
return chatapi.Conversation{}, err
}
// The requester's own words, as the first message. Only on a thread that has none: a repair
// pass must not post the opening line twice.
if err := s.openingMessage(ctx, thread, req); err != nil {
return chatapi.Conversation{}, err
}
return s.inboxRow(ctx, req.RequesterID.Int64(), thread, nil, 0)
}
// openingMessage writes the ticket's first message, and does nothing when the thread already has
// one — which is what makes the whole call safe to repeat.
func (s *Service) openingMessage(ctx context.Context, thread domain.Conversation, req chatapi.OpenTicketThreadRequest) error {
existing, err := s.repo.LastMessages(ctx, []int64{thread.ID})
if err != nil {
return fmt.Errorf("read last messages: %w", err)
}
if _, ok := existing[thread.ID]; ok {
return nil
}
if strings.TrimSpace(req.Body) == "" && len(req.Attachments) == 0 {
return nil
}
attachments := resourceKeys(req.Attachments)
if err := s.requireResources(ctx, attachments); err != nil {
return err
}
m, err := domain.NewMessage(thread.ID, req.RequesterID.Int64(), req.Body, attachments, nil, nil)
if err != nil {
return err
}
if err := s.repo.InsertMessage(ctx, &m); err != nil {
return fmt.Errorf("insert message: %w", err)
}
return nil
}
// desk is the support account's id, memoised by the account module itself.
func (s *Service) desk(ctx context.Context) (int64, error) {
support, err := s.accounts.GetSupportAccount(ctx)
if err != nil {
// Propagated as it came: the account module already named this failure, and its coded
// error is what turns a deployment with no desk row into a 500 rather than a 404 telling
// a requester their own ticket's counterparty does not exist.
return 0, err
}
return support.ID.Int64(), nil
}
func (s *Service) GetConversation(ctx context.Context, req chatapi.GetConversationRequest) (chatapi.Conversation, error) {
thread, err := s.participant(ctx, req.ActorID, req.ID)
if err != nil {
return chatapi.Conversation{}, err
}
return s.threadRow(ctx, req.ActorID.Int64(), thread)
}
// GetUnreadCount is the badge. Two numbers, because "12 unread" and "in 3 threads" are
// different things a client shows in different places.
func (s *Service) GetUnreadCount(ctx context.Context, req chatapi.UnreadCountRequest) (chatapi.UnreadCount, error) {
total, threads, err := s.repo.UnreadTotal(ctx, req.ActorID.Int64())
if err != nil {
return chatapi.UnreadCount{}, fmt.Errorf("read unread total: %w", err)
}
return chatapi.UnreadCount{Unread: total, Conversations: threads}, nil
}
// ListMessages pages a thread newest first. A thread the caller is not in is not found
// rather than forbidden — it is not theirs to know about.
func (s *Service) ListMessages(ctx context.Context, req chatapi.ListMessagesRequest) (chatapi.MessagePage, error) {
thread, err := s.participant(ctx, req.ActorID, req.ID)
if err != nil {
return chatapi.MessagePage{}, err
}
before, beforeID, err := parseCursor(req.Cursor)
if err != nil {
return chatapi.MessagePage{}, err
}
messages, err := s.repo.ListMessages(ctx, port.HistoryFilter{
ConversationID: req.ID.Int64(),
Before: before,
BeforeID: beforeID,
Limit: req.Limit + 1,
})
if err != nil {
return chatapi.MessagePage{}, fmt.Errorf("list messages: %w", err)
}
hasMore := len(messages) > req.Limit
if hasMore {
messages = messages[:req.Limit]
}
out, err := s.toAPIMessages(ctx, messages)
if err != nil {
return chatapi.MessagePage{}, err
}
page := chatapi.MessagePage{
Data: hideSupport(thread, req.ActorID.Int64(), out),
Meta: chatapi.CursorInfo{HasMore: hasMore},
}
if hasMore && len(messages) > 0 {
last := messages[len(messages)-1]
page.Meta.NextCursor = formatCursor(last.CreatedAt, last.ID)
}
return page, nil
}
// SendMessage appends to a thread the caller is in. The attachments have to be this
// module's own confirmed uploads: a message pointing at nothing is a photo that never
// renders.
func (s *Service) SendMessage(ctx context.Context, req chatapi.SendMessageRequest) (chatapi.Message, error) {
thread, err := s.participant(ctx, req.ActorID, req.ConversationID)
if err != nil {
return chatapi.Message{}, err
}
attachments := resourceKeys(req.Attachments)
if err := s.requireResources(ctx, attachments); err != nil {
return chatapi.Message{}, err
}
replyTo, err := s.replyTarget(ctx, req.ConversationID.Int64(), req.ReplyTo)
if err != nil {
return chatapi.Message{}, err
}
m, err := domain.NewMessage(req.ConversationID.Int64(), req.ActorID.Int64(), req.Body,
attachments, req.Refs, replyTo)
if err != nil {
return chatapi.Message{}, err
}
if err := s.repo.InsertMessage(ctx, &m); err != nil {
return chatapi.Message{}, fmt.Errorf("insert message: %w", err)
}
dto, err := s.toAPIMessage(ctx, m)
if err != nil {
return chatapi.Message{}, err
}
// The other side, never the sender: they already hold this row from the response they are about
// to get, and echoing it back would race their optimistic update.
if recipient := s.recipient(ctx, thread, req.ActorID.Int64()); recipient != 0 {
notify(ctx, s, recipient, MessageCreated, hideSupport(thread, recipient, []chatapi.Message{dto})[0])
}
return dto, nil
}
// replyTarget turns the reference a client sent into one the message can carry.
//
// Read rather than trusted, for one reason: the quote carries a preview of what the target
// said, so accepting a target from another conversation would read a thread the sender is
// not in out through one they are. A redacted target is allowed through — the quote then
// renders as redacted, which is the honest thing to say about an answer to it.
func (s *Service) replyTarget(ctx context.Context, conversationID int64, ref *chatapi.MessageRefRequest) (*domain.MessageRef, error) {
if ref == nil {
return nil, nil
}
target, err := s.repo.FindMessageAt(ctx, ref.ID.Int64(), ref.CreatedAt)
if err != nil {
return nil, fmt.Errorf("read reply target: %w", err)
}
if target.ConversationID != conversationID {
return nil, domain.ErrReplyOutsideThread
}
return &domain.MessageRef{ID: target.ID, CreatedAt: target.CreatedAt}, nil
}
// recipient is who a new message is pushed to, and 0 when nobody should be. On a direct thread that
// is the counterparty. On a ticket thread it is always the requester: the other side of the row is
// the desk's account, which nobody is signed in as, and staff learn about a new ticket from their
// queue rather than from a socket — so a requester writing pushes to nobody.
func (s *Service) recipient(ctx context.Context, thread domain.Conversation, writerID int64) int64 {
if !thread.Ticket() {
return thread.Other(writerID)
}
desk, err := s.desk(ctx)
if err != nil {
s.log.Error("read support account for push", "err", err)
return 0
}
if requester := thread.Counterparty(desk); requester != writerID {
return requester
}
return 0
}
// MarkConversationRead moves the caller's own mark. Absent means "everything so far",
// which is what opening a thread does.
func (s *Service) MarkConversationRead(ctx context.Context, req chatapi.MarkConversationReadRequest) (chatapi.Conversation, error) {
thread, err := s.participant(ctx, req.ActorID, req.ID)
if err != nil {
return chatapi.Conversation{}, err
}
at := time.Now()
if req.Before != nil {
at = *req.Before
}
// Staff move the desk's mark, which is what makes it shared: the next moderator opening the
// thread sees what the last one had read, and the requester's receipt says support looked.
reader := s.side(ctx, thread, req.ActorID.Int64())
if err := thread.MarkRead(reader, at); err != nil {
return chatapi.Conversation{}, err
}
if err := s.repo.SaveConversation(ctx, thread); err != nil {
return chatapi.Conversation{}, fmt.Errorf("save conversation: %w", err)
}
// The other side, never the reader: this is their own read mark advancing, which is
// nothing new to them. Sent as the side that read, so a requester is told the desk read it
// rather than which moderator did.
if other := thread.Other(reader); other != 0 {
notify(ctx, s, other, ConversationRead, chatapi.ConversationReadMark{
ConversationID: id.Of[id.Conversation](thread.ID),
ReaderID: id.Of[id.Account](reader),
ReadAt: at,
})
}
return s.threadRow(ctx, req.ActorID.Int64(), thread)
}
// UpdateMessage rewrites a body. Only the sender's own, and never a system message or one
// already unsent.
func (s *Service) UpdateMessage(ctx context.Context, req chatapi.UpdateMessageRequest) (chatapi.Message, error) {
m, err := s.repo.FindMessageAt(ctx, req.ID.Int64(), req.CreatedAt)
if err != nil {
return chatapi.Message{}, fmt.Errorf("find message: %w", err)
}
if err := m.Edit(req.ActorID.Int64(), req.Body); err != nil {
return chatapi.Message{}, err
}
if err := s.repo.SaveMessage(ctx, m); err != nil {
return chatapi.Message{}, fmt.Errorf("save message: %w", err)
}
dto, err := s.toAPIMessage(ctx, m)
if err != nil {
return chatapi.Message{}, err
}
if other := s.findOther(ctx, m.ConversationID, req.ActorID.Int64()); other != 0 {
notify(ctx, s, other, MessageUpdated, dto)
}
return dto, nil
}
// RedactMessage is unsending: the content goes, the row stays so the thread has no
// unexplained gaps.
func (s *Service) RedactMessage(ctx context.Context, req chatapi.RedactMessageRequest) error {
m, err := s.repo.FindMessageAt(ctx, req.ID.Int64(), req.CreatedAt)
if err != nil {
return fmt.Errorf("find message: %w", err)
}
if err := m.Redact(req.ActorID.Int64(), false); err != nil {
return err
}
if err := s.repo.SaveMessage(ctx, m); err != nil {
return fmt.Errorf("save message: %w", err)
}
if other := s.findOther(ctx, m.ConversationID, req.ActorID.Int64()); other != 0 {
notify(ctx, s, other, MessageDeleted, chatapi.DeletedMessageRef{
ID: id.Of[id.Message](m.ID),
ConversationID: id.Of[id.Conversation](m.ConversationID),
CreatedAt: m.CreatedAt,
})
}
return nil
}
// PostSystemMessage is another module speaking into the pair's thread — an offer card, an
// order update. It opens the thread if they have never spoken, so a caller never has to
// know whether one exists.
// PostTicketMessage posts a system message into a ticket's thread. It ensures the thread rather
// than reading it, so a ticket whose thread never opened still receives the verdict — the same
// idempotent call OpenTicketThread makes.
func (s *Service) PostTicketMessage(ctx context.Context, req chatapi.PostTicketMessageRequest) (chatapi.Message, error) {
if err := validation.AsError(s.v.Struct(req)); err != nil {
return chatapi.Message{}, err
}
desk, err := s.desk(ctx)
if err != nil {
return chatapi.Message{}, err
}
thread, err := s.repo.EnsureTicketThread(ctx, req.RequesterID.Int64(), desk, req.TicketID.Int64())
if err != nil {
return chatapi.Message{}, err
}
m, err := domain.NewSystemMessage(thread.ID, req.Body, req.Card)
if err != nil {
return chatapi.Message{}, err
}
if err := s.repo.InsertMessage(ctx, &m); err != nil {
return chatapi.Message{}, fmt.Errorf("insert message: %w", err)
}
dto, err := s.toAPIMessage(ctx, m)
if err != nil {
return chatapi.Message{}, err
}
// Pushed, unlike an offer card — that one has its own bridged event. A verdict the requester
// only sees when they refresh is the thing they are waiting for.
notify(ctx, s, req.RequesterID.Int64(), MessageCreated, dto)
return dto, nil
}
func (s *Service) PostSystemMessage(ctx context.Context, req chatapi.PostSystemMessageRequest) (chatapi.Message, error) {
if err := validation.AsError(s.v.Struct(req)); err != nil {
return chatapi.Message{}, err
}
thread, err := s.repo.EnsureConversation(ctx, req.AccountAID.Int64(), req.AccountBID.Int64())
if err != nil {
return chatapi.Message{}, fmt.Errorf("ensure conversation: %w", err)
}
m, err := domain.NewSystemMessage(thread.ID, req.Body, req.Card)
if err != nil {
return chatapi.Message{}, err
}
if err := s.repo.InsertMessage(ctx, &m); err != nil {
return chatapi.Message{}, fmt.Errorf("insert message: %w", err)
}
return s.toAPIMessage(ctx, m)
}
// GetMessage reads one message. A participant sees their own thread's; a moderator sees any,
// because a report about a message is judged on the message itself. Not found rather than
// forbidden for everyone else: a thread they are not part of is not theirs to know about.
func (s *Service) GetMessage(ctx context.Context, req chatapi.GetMessageRequest) (chatapi.Message, error) {
if err := s.v.Struct(req); err != nil {
return chatapi.Message{}, err
}
m, err := s.repo.FindMessage(ctx, req.ID.Int64())
if err != nil {
return chatapi.Message{}, fmt.Errorf("find message: %w", err)
}
thread, err := s.repo.FindConversation(ctx, m.ConversationID)
if err != nil {
return chatapi.Message{}, fmt.Errorf("find conversation: %w", err)
}
if !thread.Involves(req.ActorID.Int64()) && !s.isModerator(ctx, req.ActorID) {
return chatapi.Message{}, domain.ErrMessageNotFound
}
return s.toAPIMessage(ctx, m)
}
// hideSupport blanks the sender on a ticket thread's replies, for the requester only.
//
// The two participants of a ticket thread are the requester and the desk's own account; a moderator
// answers without being either, so "the reader is a side of this thread" is exactly "the reader is
// the requester" and needs no lookup. Staff reading their own queue keep the real sender, because a
// colleague's name is what makes a thread reviewable — the anonymity is towards the requester, whose
// complaint should be answered by the platform rather than by somebody they can go after.
func hideSupport(thread domain.Conversation, viewerID int64, messages []chatapi.Message) []chatapi.Message {
if !thread.Ticket() || !thread.Involves(viewerID) {
return messages
}
for i, m := range messages {
// A quote is a projection of a message as much as the row is: leaving its sender in
// place named the moderator who wrote the message being answered, and one public
// account read turns that id into a name.
if quote := messages[i].ReplyTo; quote != nil && quote.SenderID != nil &&
quote.SenderID.Int64() != viewerID {
quote.SenderID = nil
quote.FromSupport = true
}
if m.SenderID == nil || m.SenderID.Int64() == viewerID {
continue
}
messages[i].SenderID = nil
messages[i].FromSupport = true
}
return messages
}
// side is the account whose side of the thread the caller acts on. A moderator answering a ticket is
// not a side of the row, so every viewer-relative value — the counterparty, the read marks — would
// otherwise be computed against an account the thread has never heard of: `Counterparty` and
// `ReadMark` fall through to whichever id sorts lower, and `MarkRead` refuses outright. Staff act as
// the desk, which is also what makes the read mark shared: the next moderator inherits it.
//
// Message senders are the exception and are never mapped — see hideSupport, which anonymises only
// for a viewer the thread actually involves.
func (s *Service) side(ctx context.Context, thread domain.Conversation, actorID int64) int64 {
if thread.Involves(actorID) {
return actorID
}
desk, err := s.desk(ctx)
if err != nil {
// participant() has already refused everybody but staff on a ticket thread, so there is no
// safe answer here and the caller's own error path is the honest one.
s.log.Error("read support account for a staff view", "conversation_id", thread.ID, "err", err)
return 0
}
return desk
}
// isModerator asks the account module for the caller's role: it is a row in that module's
// table. An admin passes every moderator check.
func (s *Service) isModerator(ctx context.Context, actorID id.ID[id.Account]) bool {
me, err := s.accounts.GetMe(ctx, accountapi.GetMeRequest{ActorID: actorID})
if err != nil {
return false
}
return me.Role == accountapi.RoleModerator || me.Role == accountapi.RoleAdmin
}
// participant reads a thread the caller is in. Not found rather than forbidden: a thread
// they are not part of is not theirs to know about.
func (s *Service) participant(ctx context.Context, actorID id.ID[id.Account], conversationID id.ID[id.Conversation]) (domain.Conversation, error) {
thread, err := s.repo.FindConversation(ctx, conversationID.Int64())
if err != nil {
return domain.Conversation{}, fmt.Errorf("find conversation: %w", err)
}
if !thread.Involves(actorID.Int64()) {
// A ticket thread is answered by whichever moderator picks it up, so staff are let in
// without being a side of it. A direct thread never is: two people talking is nobody
// else's to read.
if !thread.Ticket() || !s.isModerator(ctx, actorID) {
return domain.Conversation{}, domain.ErrConversationNotFound
}
}
return thread, nil
}
// findOther loads the thread only to learn who is not the actor. UpdateMessage and
// RedactMessage authorise off the message itself, so — unlike SendMessage and
// MarkConversationRead — they hold no conversation to reuse; the lookup is as best-effort
// as the notify it feeds, since the message write it follows already committed.
func (s *Service) findOther(ctx context.Context, conversationID, actorID int64) int64 {
thread, err := s.repo.FindConversation(ctx, conversationID)
if err != nil {
s.log.Warn("realtime: find conversation for notify failed", "conversation_id", conversationID, "err", err)
return 0
}
return thread.Other(actorID)
}
// inboxRows fills a page: the last message and the unread count for every thread at once,
// and the counterparty's name per row because the account module has no batch read.
// threadRow renders a single thread, for a caller who may be staff: the unread count is asked for the
// side they act on, which is why this is not the one-element case of inboxRows — the inbox lists only
// threads the caller is a side of, a ticket read does not.
func (s *Service) threadRow(ctx context.Context, actorID int64, t domain.Conversation) (chatapi.Conversation, error) {
lastMessages, err := s.repo.LastMessages(ctx, []int64{t.ID})
if err != nil {
return chatapi.Conversation{}, fmt.Errorf("read last messages: %w", err)
}
unread, err := s.repo.UnreadCounts(ctx, s.side(ctx, t, actorID), []int64{t.ID})
if err != nil {
return chatapi.Conversation{}, fmt.Errorf("read unread counts: %w", err)
}
var last *domain.Message
if m, ok := lastMessages[t.ID]; ok {
last = &m
}
return s.inboxRow(ctx, actorID, t, last, unread[t.ID])
}
func (s *Service) inboxRows(ctx context.Context, viewerID int64, threads []domain.Conversation) ([]chatapi.Conversation, error) {
ids := make([]int64, 0, len(threads))
for _, t := range threads {
ids = append(ids, t.ID)
}
lastMessages, err := s.repo.LastMessages(ctx, ids)
if err != nil {
return nil, fmt.Errorf("read last messages: %w", err)
}
unread, err := s.repo.UnreadCounts(ctx, viewerID, ids)
if err != nil {
return nil, fmt.Errorf("read unread counts: %w", err)
}
out := make([]chatapi.Conversation, 0, len(threads))
for _, t := range threads {
var last *domain.Message
if m, ok := lastMessages[t.ID]; ok {
last = &m
}
row, err := s.inboxRow(ctx, viewerID, t, last, unread[t.ID])
if err != nil {
return nil, err
}
out = append(out, row)
}
return out, nil
}
// inboxRow renders one thread as its viewer sees it. `viewerID` is the caller and `side` is the side
// of the row they act on — the same account except for staff on a ticket thread, where the two differ
// on purpose: the counterparty and the read marks are the desk's, while the last message is
// anonymised for the requester only.
func (s *Service) inboxRow(ctx context.Context, viewerID int64, t domain.Conversation, last *domain.Message, unread int64) (chatapi.Conversation, error) {
side := s.side(ctx, t, viewerID)
counterparty, err := s.accounts.GetPublicAccount(ctx, accountapi.GetPublicAccountRequest{
ID: id.Of[id.Account](t.Counterparty(side)),
})
if err != nil {
return chatapi.Conversation{}, fmt.Errorf("read counterparty: %w", err)
}
row := chatapi.Conversation{
ID: id.Of[id.Conversation](t.ID),
Counterparty: accountapi.AccountSummary{
ID: counterparty.ID, Name: counterparty.Name, Avatar: counterparty.Avatar,
},
LastMessageAt: t.LastMessageAt,
Unread: unread,
ReadAt: t.ReadMark(side),
CounterpartyReadAt: t.CounterpartyReadMark(side),
CreatedAt: t.CreatedAt,
}
// A ticket thread is answered in the support screen rather than the inbox, so the client has to
// be able to tell them apart from the row alone.
if t.TicketID != nil {
row.TicketID = new(id.Of[id.Ticket](*t.TicketID))
}
if last != nil {
message, err := s.toAPIMessage(ctx, *last)
if err != nil {
return chatapi.Conversation{}, err
}
// Anonymised here too, not only in the thread: this row is what `GET /conversations` renders,
// and a public account read turns a moderator's id into their name and their shop.
row.LastMessage = &hideSupport(t, viewerID, []chatapi.Message{message})[0]
}
return row, nil
}
func (s *Service) toAPIMessages(ctx context.Context, messages []domain.Message) ([]chatapi.Message, error) {
// One resource read for the whole page rather than one per message.
var keys []int64
for _, m := range messages {
keys = append(keys, m.Attachments...)
}
images, err := s.resources(ctx, keys)
if err != nil {
return nil, err
}
quotes, err := s.quotes(ctx, messages)
if err != nil {
return nil, err
}
out := make([]chatapi.Message, 0, len(messages))
for _, m := range messages {
out = append(out, buildMessage(m, images, quotes))
}
return out, nil
}
func (s *Service) toAPIMessage(ctx context.Context, m domain.Message) (chatapi.Message, error) {
images, err := s.resources(ctx, m.Attachments)
if err != nil {
return chatapi.Message{}, err
}
quotes, err := s.quotes(ctx, []domain.Message{m})
if err != nil {
return chatapi.Message{}, err
}
return buildMessage(m, images, quotes), nil
}
// quotes resolves what a page of replies points at, keyed by the quoted message's id — one
// read for the page. Live rather than snapshotted at send time, so an edit to the original
// shows through and a redaction reads as one.
func (s *Service) quotes(ctx context.Context, messages []domain.Message) (map[int64]chatapi.MessageQuote, error) {
var refs []domain.MessageRef
for _, m := range messages {
if m.ReplyTo != nil {
refs = append(refs, *m.ReplyTo)
}
}
if len(refs) == 0 {
return nil, nil
}
found, err := s.repo.QuotedMessages(ctx, refs)
if err != nil {
return nil, fmt.Errorf("read quoted messages: %w", err)
}
out := make(map[int64]chatapi.MessageQuote, len(found))
for key, quoted := range found {
out[key] = buildQuote(quoted)
}
return out, nil
}
// previewRunes caps a quote at what a bubble can hold. Runes rather than bytes: a
// Vietnamese sentence is mostly multi-byte, and cutting on a byte boundary produces a
// replacement character instead of a shorter sentence.
const previewRunes = 120
func buildQuote(m domain.Message) chatapi.MessageQuote {
out := chatapi.MessageQuote{
ID: id.Of[id.Message](m.ID),
CreatedAt: m.CreatedAt,
Preview: truncateRunes(m.Body, previewRunes),
Attachments: len(m.Attachments),
Redacted: !m.IsLive(),
}
if m.SenderID != 0 {
out.SenderID = new(id.Of[id.Account](m.SenderID))
}
return out
}
func truncateRunes(s string, limit int) string {
runes := []rune(s)
if len(runes) <= limit {
return s
}
return string(runes[:limit]) + "…"
}
func buildMessage(m domain.Message, images map[int64]common.ResourceDTO, quotes map[int64]chatapi.MessageQuote) chatapi.Message {
out := chatapi.Message{
ID: id.Of[id.Message](m.ID),
ConversationID: id.Of[id.Conversation](m.ConversationID),
Type: m.Type,
Body: m.Body,
Images: pick(images, m.Attachments),
Refs: m.Refs,
Card: m.Card,
CreatedAt: m.CreatedAt,
EditedAt: m.EditedAt,
DeletedAt: m.DeletedAt,
}
// Null on a system message: that one is the backend's word, not a person's.
if m.SenderID != 0 {
out.SenderID = new(id.Of[id.Account](m.SenderID))
}
if m.ReplyTo != nil {
if quote, ok := quotes[m.ReplyTo.ID]; ok {
out.ReplyTo = "e
} else {
// The reference outlived its row, which a deleted conversation is the only way
// to arrange. Naming it as unavailable beats dropping the quote, which would
// render an answer as though it had answered nothing.
out.ReplyTo = &chatapi.MessageQuote{
ID: id.Of[id.Message](m.ReplyTo.ID),
CreatedAt: m.ReplyTo.CreatedAt,
Redacted: true,
}
}
}
return out
}
// requireResources refuses an attachment that names no confirmed upload of this module's.
func (s *Service) requireResources(ctx context.Context, keys []int64) error {
if len(keys) == 0 {
return nil
}
found, err := s.resources(ctx, keys)
if err != nil {
return err
}
for _, key := range keys {
if _, ok := found[key]; !ok {
return domain.ErrAttachmentNotFound
}
}
return nil
}
func (s *Service) resources(ctx context.Context, keys []int64) (map[int64]common.ResourceDTO, error) {
if len(keys) == 0 {
return map[int64]common.ResourceDTO{}, nil
}
return s.uploads.Resolve(ctx, keys)
}
// pick keeps the sender's order, which is the only thing the array encodes.
func pick(found map[int64]common.ResourceDTO, keys []int64) []common.ResourceDTO {
out := make([]common.ResourceDTO, 0, len(keys))
for _, key := range keys {
if res, ok := found[key]; ok {
out = append(out, res)
}
}
return out
}
func resourceKeys(ids []id.ID[id.Resource]) []int64 {
out := make([]int64, 0, len(ids))
for _, rid := range ids {
out = append(out, rid.Int64())
}
return out
}
// The cursor is a (timestamp, id) tuple; common owns the format, and the reason it is a tuple.
func formatCursor(at time.Time, id int64) string {
return common.FormatCursor(at.UnixNano(), id)
}
func parseCursor(cursor string) (time.Time, int64, error) {
nanos, id, err := common.ParseCursor(cursor)
if err != nil || id == 0 {
return time.Time{}, 0, err
}
return time.Unix(0, nanos), id, nil
}