-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathSlackDailyDigestService.java
More file actions
3226 lines (2943 loc) · 159 KB
/
Copy pathSlackDailyDigestService.java
File metadata and controls
3226 lines (2943 loc) · 159 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
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
package com.dbaagent.service;
import com.dbaagent.dto.SlowQueryHistorySummary;
import com.dbaagent.dto.SlowQueryInsightsResponse;
import com.dbaagent.model.CapacityForecast;
import com.dbaagent.model.ColumnInfo;
import com.dbaagent.model.ConnectionAccessGrant;
import com.dbaagent.model.ConnectionRequest;
import com.dbaagent.model.DatabaseEvent;
import com.dbaagent.model.DatabaseObject;
import com.dbaagent.model.DigestDeliveryMethod;
import com.dbaagent.model.GrowthAnomaly;
import com.dbaagent.model.IndexRecommendationEntity;
import com.dbaagent.model.IndexRecommendationEvidence;
import com.dbaagent.model.LockContention;
import com.dbaagent.model.PerformanceAction;
import com.dbaagent.model.PerformanceSnapshot;
import com.dbaagent.model.PersonaTag;
import com.dbaagent.model.QueryFingerprint;
import com.dbaagent.util.QueryLabeler;
import com.dbaagent.model.QueryRequest;
import com.dbaagent.model.QueryResult;
import com.dbaagent.model.Role;
import com.dbaagent.model.SchemaChange;
import com.dbaagent.model.SchemaSnapshot;
import com.dbaagent.model.SlackChannelBinding;
import com.dbaagent.model.SlackDigestLog;
import com.dbaagent.model.SlackUserLink;
import com.dbaagent.model.SlowQuery;
import com.dbaagent.model.SlowQueryAnalysis;
import com.dbaagent.model.TableStatsHistory;
import com.dbaagent.model.User;
import com.dbaagent.model.UserDigestPreference;
import com.dbaagent.model.digest.DigestAssemblyResult;
import com.dbaagent.model.digest.DigestInsight;
import com.dbaagent.repository.AuthLoginChallengeRepository;
import com.dbaagent.repository.UserDigestPreferenceRepository;
import com.dbaagent.repository.UserRepository;
import com.dbaagent.repository.CapacityForecastRepository;
import com.dbaagent.repository.ConnectionAccessGrantRepository;
import com.dbaagent.repository.DatabaseEventRepository;
import com.dbaagent.repository.GrowthAnomalyRepository;
import com.dbaagent.repository.LockContentionRepository;
import com.dbaagent.repository.QueryFingerprintRepository;
import com.dbaagent.repository.SchemaChangeRepository;
import com.dbaagent.repository.SlackChannelBindingRepository;
import com.dbaagent.repository.SlackDigestLogRepository;
import com.dbaagent.repository.TableStatsHistoryRepository;
import com.dbaagent.service.digest.DigestCronMatcher;
import com.dbaagent.service.digest.DigestInsightAssemblerService;
import com.slack.api.Slack;
import com.slack.api.methods.MethodsClient;
import com.slack.api.methods.SlackApiException;
import com.slack.api.methods.request.chat.ChatPostMessageRequest;
import com.slack.api.methods.request.conversations.ConversationsOpenRequest;
import com.slack.api.methods.response.conversations.ConversationsOpenResponse;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.lang.Nullable;
import org.springframework.stereotype.Service;
import java.io.IOException;
import java.time.LocalDateTime;
import java.time.ZoneId;
import java.time.Instant;
import java.time.Duration;
import java.time.format.DateTimeFormatter;
import java.util.ArrayList;
import java.util.Comparator;
import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Locale;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.Set;
import java.util.stream.Collectors;
import org.springframework.beans.factory.annotation.Value;
@Service
@RequiredArgsConstructor
@Slf4j
public class SlackDailyDigestService {
/**
* When true (default), the daily digest is delivered only to Slack channel
* bindings owned by a DeepSQL admin. A binding is created per channel/DM when
* someone picks a connection, so without this filter every linked user gets
* the digest DM'd. Set false to restore broadcast-to-all delivery.
*/
@Value("${slack.digest.admins-only:true}")
private boolean digestAdminsOnly;
/**
* Global / legacy digest schedule. Also the default when a preference leaves
* {@code cronExpression} blank. Interpreted in UTC for legacy; per-user prefs
* use their own timezone.
*/
@Value("${slack.daily-digest.cron:0 0 9 * * *}")
private String globalDigestCron;
private final SlackRuntimeSettingsService slackRuntimeSettingsService;
private final SlackChannelBindingRepository channelBindingRepository;
private final CredentialService credentialService;
private final ConnectionService connectionService;
private final PerformanceInsightsService performanceInsightsService;
private final SlowQueryService slowQueryService;
private final SlowQueryHistoryService slowQueryHistoryService;
private final SlowQueryInsightsService slowQueryInsightsService;
private final SlowQueryAnalyticsService slowQueryAnalyticsService;
private final PerformanceActionAggregatorService actionAggregatorService;
private final EnhancedSqlParserService sqlParserService;
private final QueryExecutorService queryExecutorService;
private final TableGrowthMonitoringService tableGrowthMonitoringService;
private final SchemaChangeTrackingService schemaChangeTrackingService;
private final TableStatsHistoryRepository tableStatsHistoryRepository;
private final GrowthAnomalyRepository growthAnomalyRepository;
private final CapacityForecastRepository capacityForecastRepository;
private final SchemaChangeRepository schemaChangeRepository;
private final SlackDigestLogRepository digestLogRepository;
private final SlackUserLinkService slackUserLinkService;
private final LockContentionRepository lockContentionRepository;
private final QueryFingerprintRepository queryFingerprintRepository;
private final DatabaseEventRepository databaseEventRepository;
private final ConnectionAccessGrantRepository connectionAccessGrantRepository;
private final AuthLoginChallengeRepository authLoginChallengeRepository;
private final IndexAdvisorService indexAdvisorService;
private final IndexRecommendationService indexRecommendationService;
private final UserDigestPreferenceRepository userDigestPreferenceRepository;
private final DigestInsightAssemblerService digestInsightAssemblerService;
private final UserRepository userRepository;
private static final DateTimeFormatter DATE_FMT = DateTimeFormatter.ofPattern("MMM d, yyyy");
private static final int TOP_TABLES = 5;
private static final int TOP_SLOW_QUERIES = 3;
private static final List<String> DIGEST_TITLES = List.of(
"🗄️ DB Health Briefing",
"🗄️ Database Situation Report",
"🗄️ Engineering DB Readout"
);
private static final List<String> SUMMARY_TITLES = List.of(
"🧭 EXECUTIVE READ",
"🧭 WHAT CHANGED",
"🧭 FIRST TAKE"
);
private static final List<String> SPOTLIGHT_TITLES = List.of(
"🔦 TODAY'S SPOTLIGHT",
"🔦 LEAD SIGNAL",
"🔦 WHAT DESERVES ATTENTION"
);
private static final List<String> SLOW_TITLES = List.of(
"🐢 SLOW QUERY SNAPSHOT",
"🐢 QUERY PRESSURE BOARD",
"🐢 WORKLOAD RISK CHECK"
);
private static final List<String> CUSTOMERS_TITLES = List.of(
"👥 TOP AFFECTED CUSTOMERS",
"👥 CUSTOMERS BEARING THE LOAD",
"👥 SLOWNESS BY CUSTOMER"
);
private static final List<String> FREE_WIN_TITLES = List.of(
"🎯 FREE WIN",
"🎯 LOW-EFFORT LEVER",
"🎯 HIGHEST-ROI MOVE"
);
private static final List<String> GROWTH_TITLES = List.of(
"📈 GROWTH WATCH",
"📈 CAPACITY WATCH",
"📈 FOOTPRINT WATCH"
);
private static final List<String> SCHEMA_TITLES = List.of(
"🔧 SCHEMA EVENTS (last 24h)",
"🔧 SCHEMA DELTAS (last 24h)",
"🔧 STRUCTURE CHANGES (last 24h)"
);
private static final List<String> SECURITY_TITLES = List.of(
"🛡️ SECURITY & PRIVILEGED ACTIONS (last 24h)",
"🛡️ ACCESS WATCH (last 24h)",
"🛡️ TRUST BOUNDARY (last 24h)"
);
private static final List<String> CONCURRENCY_TITLES = List.of(
"🔒 CONCURRENCY HOTSPOTS (last 24h)",
"🔒 LOCK & WAIT REPORT (last 24h)",
"🔒 BLOCKING ACTIVITY (last 24h)"
);
private static final List<String> WASTE_TITLES = List.of(
"💸 SILENT WASTE",
"💸 RECLAIMABLE STORAGE",
"💸 HIDDEN COSTS"
);
private static final List<String> NEWCOMER_TITLES = List.of(
"🆕 NEWCOMER QUERIES (last 24h)",
"🆕 NEW HEAVY HITTERS",
"🆕 RISING WORKLOAD"
);
private static final List<String> SAVINGS_TITLES = List.of(
"💰 POTENTIAL SAVINGS",
"💰 MONEY ON THE TABLE",
"💰 COST RECOVERY OPPORTUNITY"
);
private static final List<String> INDEX_WINS_TITLES = List.of(
"📈 INDEX WINS",
"📈 WORKLOAD-WEIGHTED INDEX ADVISOR",
"📈 INDEXES TO ADD (RANKED BY ROI)"
);
private static final List<String> DEEP_DIVE_TITLES = List.of(
"🔎 TODAY'S DEEP DIVE",
"🔎 ZOOM IN",
"🔎 ONE THING WORTH KNOWING"
);
private static final List<String> FOOTERS = List.of(
"Powered by DeepSQL · reply with any question about your DB",
"Powered by DeepSQL · ask follow-up questions right here",
"Powered by DeepSQL · use this thread to dig deeper"
);
// "savings" leads every rotation — money is the universal hook. "deep_dive" is a rotating
// angle picked daily so the digest doesn't read the same two days in a row. The rest of the
// rotation preserves prior section ordering so the macro briefing keeps its primacy.
// "customers" sits right after "slow" in every rotation — they're the
// same data sliced along a different axis (queries × customers instead
// of queries × time), so they read naturally adjacent.
private static final List<String> SECTION_ORDER_A = List.of(
"savings", "index_wins", "summary", "spotlight", "slow", "customers", "concurrency", "newcomers", "free", "growth", "waste", "schema", "security", "deep_dive");
private static final List<String> SECTION_ORDER_B = List.of(
"savings", "index_wins", "summary", "spotlight", "concurrency", "growth", "slow", "customers", "deep_dive", "newcomers", "schema", "security", "free", "waste");
private static final List<String> SECTION_ORDER_C = List.of(
"savings", "index_wins", "summary", "slow", "customers", "deep_dive", "spotlight", "newcomers", "concurrency", "schema", "security", "growth", "waste", "free");
private static final String CONNECTION_UNAVAILABLE_HEADLINE = "live database access unavailable";
private static final List<String> LABEL_COLUMN_CANDIDATES = List.of(
"name", "title", "display_name", "displayname", "label", "full_name", "customer_name",
"shop_name", "property_name", "account_name", "customer_name", "slug", "code"
);
private record DigestStyle(
String digestTitle,
String summaryTitle,
String spotlightTitle,
String slowTitle,
String customersTitle,
String freeWinTitle,
String growthTitle,
String schemaTitle,
String securityTitle,
String concurrencyTitle,
String wasteTitle,
String newcomerTitle,
String savingsTitle,
String indexWinsTitle,
String deepDiveTitle,
String footer,
List<String> order,
long digestCount
) {}
private record GrowthDigestData(
List<TableStatsHistory> latestSnapshots,
List<TableStatsHistory> topGrowing,
Map<String, CapacityForecast> forecastByTable,
Map<String, GrowthAnomaly> anomalyByTable,
boolean bootstrapped
) {}
private record SchemaDigestData(
List<SchemaChange> changes,
boolean comparedFromSnapshots,
@Nullable LocalDateTime currentSnapshotAt,
@Nullable LocalDateTime previousSnapshotAt
) {}
public record TriggerResult(boolean triggered, String message) {}
/**
* Runs the hybrid digest in a virtual thread so the HTTP trigger returns immediately.
* Uses per-user personalized delivery when prefs exist; otherwise legacy broadcast.
*/
public void sendDailyDigestAsync() {
Thread.ofVirtual().name("digest-send").start(this::sendDailyDigestHybrid);
}
public TriggerResult triggerDigest() {
List<String> connectionIds = digestConnectionIds();
if (connectionIds.isEmpty()) {
return new TriggerResult(false, "No database connections available for digest generation.");
}
boolean slackEnabled = isSlackDeliveryEnabled();
boolean perUserMode = isPerUserModeEnabled();
long bindingCount = channelBindingRepository.count();
sendDailyDigestAsync();
if (!slackEnabled) {
return new TriggerResult(true, "Digest generation started. Slack delivery is disabled in this environment.");
}
if (perUserMode) {
return new TriggerResult(true,
"Personalized digest delivery started (Slack DMs per user prefs when linked). Channel bindings are not required for DMs.");
}
if (bindingCount == 0) {
return new TriggerResult(true, "Digest generation started. No Slack channel bindings found, so it will only appear in the app.");
}
return new TriggerResult(true, "Digest generation and Slack delivery are running in the background.");
}
public void sendDailyDigest() {
List<String> connectionIds = digestConnectionIds();
if (connectionIds.isEmpty()) {
log.info("No connections available — skipping daily digest generation");
return;
}
List<SlackChannelBinding> bindings = channelBindingRepository.findAll();
Map<String, List<SlackChannelBinding>> bindingsByConnection = bindings.stream()
.filter(binding -> binding.getDefaultConnectionId() != null && !binding.getDefaultConnectionId().isBlank())
.collect(Collectors.groupingBy(SlackChannelBinding::getDefaultConnectionId));
boolean slackEnabled = isSlackDeliveryEnabled();
log.info(
"Generating daily digest for {} connection(s); Slack delivery enabled={}, bound channels={}",
connectionIds.size(),
slackEnabled,
bindings.size()
);
for (String connectionId : connectionIds) {
SlackDigestLog logEntry = new SlackDigestLog();
logEntry.setConnectionId(connectionId);
logEntry.setConnectionName(connectionName(connectionId));
logEntry.setChannelId(null);
logEntry.setSentAt(LocalDateTime.now());
try {
String message = buildRichDigest(connectionId);
logEntry.setContent(message);
logEntry.setHeadline(extractHeadline(message));
List<SlackChannelBinding> allBindings = bindingsByConnection.getOrDefault(connectionId, List.of());
List<SlackChannelBinding> connectionBindings = filterAdminRecipients(allBindings);
if (digestAdminsOnly && connectionBindings.size() < allBindings.size()) {
log.info("Digest for connection {}: delivering to {}/{} admin-owned binding(s) (admins-only filter)",
connectionId, connectionBindings.size(), allBindings.size());
}
if (!slackEnabled || connectionBindings.isEmpty()) {
logEntry.setStatus("GENERATED");
if (!slackEnabled) {
logEntry.setErrorMessage("Slack delivery disabled");
} else {
logEntry.setErrorMessage("No Slack channel bindings");
}
log.info(
"Digest generated for connection {} without Slack delivery (enabled={}, bindings={})",
connectionId,
slackEnabled,
connectionBindings.size()
);
} else {
List<String> failures = sendToBoundChannels(connectionBindings, message);
if (failures.isEmpty()) {
logEntry.setStatus("SENT");
log.info(
"Daily digest generated and sent to {} channel(s) for connection {}",
connectionBindings.size(),
connectionId
);
} else if (failures.size() < connectionBindings.size()) {
logEntry.setStatus("PARTIAL");
logEntry.setErrorMessage(String.join(" | ", failures));
} else {
logEntry.setStatus("FAILED");
logEntry.setErrorMessage(String.join(" | ", failures));
}
}
} catch (Exception e) {
logEntry.setStatus("FAILED");
logEntry.setErrorMessage(e.getMessage());
log.error("Failed to generate daily digest for connection {}: {}",
connectionId, e.getMessage(), e);
} finally {
try { digestLogRepository.save(logEntry); } catch (Exception ex) {
log.warn("Could not persist digest log: {}", ex.getMessage());
}
}
}
}
private boolean isSlackDeliveryEnabled() {
SlackRuntimeSettingsService.SlackRuntimeConfig config = slackRuntimeSettingsService.current();
return config.enabled() && config.botToken() != null && !config.botToken().isBlank();
}
/**
* Check if the system has per-user digest preferences configured.
* When true, digest delivery uses UserDigestPreference; when false, uses legacy
* channel-broadcast mode.
*
* <p>This is the foundation for role-aware digests (PR2). In PR1, this method
* allows the service to detect whether per-user mode is active, though the
* actual per-user delivery logic comes later.
*/
public boolean isPerUserModeEnabled() {
return userDigestPreferenceRepository.hasAnyPreferences();
}
/**
* Get all enabled digest preferences for a connection.
* Returns users who should receive a digest for this connection based on their
* UserDigestPreference settings.
*
* <p>For PR1, this returns the preferences but the actual personalized delivery
* is implemented in PR2. The current digest flow continues using the legacy
* channel-broadcast approach.
*
* @param connectionId the connection to get recipients for
* @return list of enabled preferences, or empty if no per-user preferences exist
*/
public List<UserDigestPreference> getDigestRecipients(String connectionId) {
return userDigestPreferenceRepository.findEnabledForConnection(connectionId);
}
/**
* Get digest mode summary for logging and debugging.
*/
public DigestModeInfo getDigestModeInfo() {
boolean perUserMode = isPerUserModeEnabled();
long preferenceCount = perUserMode ? userDigestPreferenceRepository.countByEnabledTrue() : 0;
List<String> users = perUserMode
? userDigestPreferenceRepository.findDistinctUsernamesWithEnabledPreferences()
: List.of();
return new DigestModeInfo(perUserMode, preferenceCount, users.size());
}
public record DigestModeInfo(boolean perUserMode, long enabledPreferences, int distinctUsers) {}
private List<String> digestConnectionIds() {
Set<String> ids = new HashSet<>();
credentialService.getAllConnections().stream()
.map(c -> c.getId())
.filter(id -> id != null && !id.isBlank())
.forEach(ids::add);
channelBindingRepository.findAll().stream()
.map(SlackChannelBinding::getDefaultConnectionId)
.filter(id -> id != null && !id.isBlank())
.forEach(ids::add);
return ids.stream().sorted().toList();
}
/**
* Keep only bindings owned by a DeepSQL admin when {@code slack.digest.admins-only}
* is on. The binding's {@code updatedBy} is the DeepSQL username that bound the
* channel/DM; a non-admin or unknown owner is excluded so the digest isn't
* broadcast to every linked user.
*/
private List<SlackChannelBinding> filterAdminRecipients(List<SlackChannelBinding> bindings) {
if (!digestAdminsOnly) {
return bindings;
}
return bindings.stream().filter(this::isAdminOwned).toList();
}
private boolean isAdminOwned(SlackChannelBinding binding) {
// updated_by is the Slack user id that bound the channel/DM — resolve it to
// the linked DeepSQL user and check their role.
String slackUserId = binding.getUpdatedBy();
if (slackUserId == null || slackUserId.isBlank()) {
return false;
}
return slackUserLinkService.resolveLinkedUser(binding.getTeamId(), slackUserId)
.map(SlackUserLinkService.LinkedUser::admin)
.orElse(false);
}
private List<String> sendToBoundChannels(List<SlackChannelBinding> bindings, String message) {
return bindings.stream()
.map(binding -> {
try {
postMessage(binding.getChannelId(), message);
return null;
} catch (Exception e) {
log.error("Failed to send daily digest to channel {}: {}",
binding.getChannelId(), e.getMessage(), e);
return binding.getChannelId() + ": " + e.getMessage();
}
})
.filter(error -> error != null && !error.isBlank())
.toList();
}
// ─────────────────────────────────────────────────────────────────────────
// Per-recipient personalized digest delivery (PR3)
// ─────────────────────────────────────────────────────────────────────────
/**
* Send personalized digests to all users with enabled preferences for a connection.
* Uses DigestAssemblyResult to generate role-aware, persona-tailored content.
*
* <p>Two users with different personas on the same connection get different
* Slack digests for the same time window.
*
* @param connectionId the connection to generate digests for
* @return summary of delivery results
*/
public PersonalizedDeliveryResult sendPersonalizedDigests(String connectionId) {
List<UserDigestPreference> recipients = userDigestPreferenceRepository
.findEnabledForConnection(connectionId);
if (recipients.isEmpty()) {
log.debug("No per-user preferences for connection {}; skipping personalized delivery", connectionId);
return new PersonalizedDeliveryResult(0, 0, 0, List.of());
}
boolean slackEnabled = isSlackDeliveryEnabled();
if (!slackEnabled) {
log.info("Slack delivery disabled; generating {} personalized digest(s) without sending", recipients.size());
}
String connName = connectionName(connectionId);
LocalDateTime since = getLastDigestTime(connectionId);
List<String> errors = new ArrayList<>();
int sent = 0;
int generated = 0;
for (UserDigestPreference pref : recipients) {
if (pref.getDeliveryMethod() != DigestDeliveryMethod.SLACK_DM) {
continue; // Only Slack DM for now; EMAIL/WhatsApp in PR4
}
try {
PersonalizedDigestResult result = sendPersonalizedDigestToUser(
connectionId, connName, pref, since, slackEnabled);
if (result.generated()) {
generated++;
if (result.sent()) {
sent++;
}
}
if (result.error() != null) {
errors.add(pref.getUsername() + ": " + result.error());
}
} catch (Exception e) {
log.error("Failed to deliver personalized digest to user {} for connection {}: {}",
pref.getUsername(), connectionId, e.getMessage(), e);
errors.add(pref.getUsername() + ": " + e.getMessage());
}
}
log.info("Personalized digest delivery for connection {}: {} generated, {} sent, {} errors",
connectionId, generated, sent, errors.size());
return new PersonalizedDeliveryResult(recipients.size(), generated, sent, errors);
}
/**
* Send a personalized digest to a single user based on their preference.
*/
private PersonalizedDigestResult sendPersonalizedDigestToUser(
String connectionId, String connName,
UserDigestPreference pref, LocalDateTime since,
boolean slackEnabled) {
String username = pref.getUsername();
PersonaTag personaTag = pref.getPersonaTag();
// Resolve user's built-in role. Custom roles have no Role enum — rank as
// DEVELOPER rather than NPE on role.name(). Persist the stored role code
// on the log so the audit trail still names ANALYST, not a silent fallback.
User recipient = userRepository.findByUsernameIgnoreCase(username).orElse(null);
Role role = recipient != null && recipient.getRoleEnum() != null
? recipient.getRoleEnum()
: Role.DEVELOPER;
String roleCode = recipient != null ? recipient.getRoleCode() : Role.DEVELOPER.name();
// Assemble personalized digest
DigestAssemblyResult assembly = digestInsightAssemblerService.assembleDigest(
username, connectionId, role, personaTag, since);
// Format digest message
String message = formatPersonalizedDigest(assembly, connName, personaTag);
// Create log entry
SlackDigestLog logEntry = new SlackDigestLog();
logEntry.setConnectionId(connectionId);
logEntry.setConnectionName(connName);
logEntry.setSentAt(LocalDateTime.now());
logEntry.setContent(message);
logEntry.setHeadline(assembly.getHeadline());
logEntry.setRecipientUsername(username);
logEntry.setRecipientRole(roleCode);
logEntry.setPersonaTag(personaTag);
logEntry.setDeliveryMethod(DigestDeliveryMethod.SLACK_DM);
logEntry.setPreferenceId(pref.getId());
logEntry.setPersonalized(true);
boolean sent = false;
String error = null;
if (slackEnabled && !assembly.isEmpty()) {
try {
String dmChannelId = openDmChannel(username);
if (dmChannelId != null) {
postMessage(dmChannelId, message);
logEntry.setChannelId(dmChannelId);
logEntry.setStatus("SENT");
sent = true;
log.debug("Sent personalized digest to {} (persona={}, {} insights)",
username, personaTag, assembly.getInsights().size());
} else {
logEntry.setStatus("FAILED");
error = "Could not open DM channel - user not linked to Slack";
logEntry.setErrorMessage(error);
}
} catch (Exception e) {
logEntry.setStatus("FAILED");
error = e.getMessage();
logEntry.setErrorMessage(error);
log.error("Failed to send personalized digest DM to {}: {}", username, e.getMessage());
}
} else {
logEntry.setStatus("GENERATED");
if (!slackEnabled) {
logEntry.setErrorMessage("Slack delivery disabled");
} else if (assembly.isEmpty()) {
logEntry.setErrorMessage("No insights to deliver");
}
}
// Persist log
try {
digestLogRepository.save(logEntry);
} catch (Exception e) {
log.warn("Could not persist personalized digest log for {}: {}", username, e.getMessage());
}
return new PersonalizedDigestResult(true, sent, error);
}
/**
* Open a DM channel with a user via their linked Slack account.
* Returns the channel ID for sending messages, or null if user not linked.
*/
private String openDmChannel(String deepsqlUsername) {
List<SlackUserLink> links = slackUserLinkService.getLinkedSlackAccounts(deepsqlUsername);
if (links.isEmpty()) {
log.debug("User {} has no linked Slack accounts for DM delivery", deepsqlUsername);
return null;
}
// Use the first linked account (most users have one)
SlackUserLink link = links.get(0);
String botToken = slackRuntimeSettingsService.current().botToken();
if (botToken == null || botToken.isBlank()) {
log.warn("No Slack bot token configured; cannot open DM channel");
return null;
}
MethodsClient client = Slack.getInstance().methods(botToken);
try {
ConversationsOpenResponse response = client.conversationsOpen(
ConversationsOpenRequest.builder()
.users(List.of(link.getSlackUserId()))
.build());
if (response.isOk() && response.getChannel() != null) {
return response.getChannel().getId();
} else {
log.error("Failed to open DM channel with Slack user {}: {}",
link.getSlackUserId(), response.getError());
return null;
}
} catch (IOException | SlackApiException e) {
log.error("Failed to open DM channel with Slack user {}: {}",
link.getSlackUserId(), e.getMessage());
return null;
}
}
/**
* Format a personalized digest message from DigestAssemblyResult.
* EXEC personas get tight 3-bullet executive summaries.
*/
private String formatPersonalizedDigest(DigestAssemblyResult assembly, String connName, PersonaTag personaTag) {
StringBuilder sb = new StringBuilder();
DigestStyle style = chooseStyle(assembly.getConnectionId());
// Header
sb.append("*").append(style.digestTitle()).append(" — ").append(connName).append("*\n");
sb.append("_").append(LocalDateTime.now().format(DATE_FMT))
.append(" · ").append(assembly.getHeadline()).append("_\n");
if (personaTag != null) {
sb.append("_Personalized for: ").append(personaTag.getDisplayName()).append("_\n");
}
sb.append("────────────────────────\n\n");
if (assembly.isEmpty()) {
sb.append("✓ No new insights since your last digest — your database is running smoothly.\n\n");
sb.append("_").append(style.footer()).append("_");
return sb.toString();
}
// EXEC persona: tight executive summary (3 bullets max)
if (personaTag == PersonaTag.EXEC) {
formatExecDigest(sb, assembly, style);
} else {
formatStandardDigest(sb, assembly, style);
}
sb.append("\n_").append(style.footer()).append("_");
return sb.toString();
}
/**
* Format EXEC digest: 3 bullets + decision ask. Keep it tight.
*/
private void formatExecDigest(StringBuilder sb, DigestAssemblyResult assembly, DigestStyle style) {
sb.append("*🧭 EXECUTIVE SUMMARY*\n");
List<String> execSummary = assembly.getExecutiveSummary();
if (execSummary != null && !execSummary.isEmpty()) {
for (String bullet : execSummary.stream().limit(3).toList()) {
sb.append("• ").append(bullet).append("\n");
}
} else {
// Generate summary from top insights
List<DigestInsight> top = assembly.getTopInsights(3);
for (DigestInsight insight : top) {
String emoji = getInsightEmoji(insight);
sb.append("• ").append(emoji).append(" ").append(insight.getHeadline()).append("\n");
}
}
sb.append("\n");
// Decision ask
String decisionAsk = assembly.getDecisionAsk();
if (decisionAsk != null && !decisionAsk.isBlank()) {
sb.append("*📋 ACTION NEEDED*\n");
sb.append(decisionAsk).append("\n\n");
}
// Stats line
sb.append("_").append(assembly.getInsights().size()).append(" insight");
if (assembly.getInsights().size() != 1) sb.append("s");
sb.append(" available for detailed review_\n");
}
/**
* Format standard digest for non-EXEC personas.
*/
private void formatStandardDigest(StringBuilder sb, DigestAssemblyResult assembly, DigestStyle style) {
List<DigestInsight> insights = assembly.getInsights();
// Critical insights first
List<DigestInsight> critical = insights.stream()
.filter(i -> i.getSeverity() >= 90)
.toList();
if (!critical.isEmpty()) {
sb.append("*⚠️ CRITICAL*\n");
for (DigestInsight insight : critical.stream().limit(3).toList()) {
formatInsight(sb, insight);
}
sb.append("\n");
}
// High-priority insights
List<DigestInsight> high = insights.stream()
.filter(i -> i.getSeverity() >= 70 && i.getSeverity() < 90)
.toList();
if (!high.isEmpty()) {
sb.append("*🔶 HIGH PRIORITY*\n");
for (DigestInsight insight : high.stream().limit(3).toList()) {
formatInsight(sb, insight);
}
sb.append("\n");
}
// Other insights summary
List<DigestInsight> other = insights.stream()
.filter(i -> i.getSeverity() < 70)
.toList();
if (!other.isEmpty()) {
sb.append("*📋 OTHER INSIGHTS*\n");
for (DigestInsight insight : other.stream().limit(5).toList()) {
formatInsight(sb, insight);
}
if (other.size() > 5) {
sb.append("_... and ").append(other.size() - 5).append(" more_\n");
}
sb.append("\n");
}
// Suppressed/filtered counts
if (assembly.getSuppressedDuplicates() > 0 || assembly.getFilteredAcknowledged() > 0) {
sb.append("_");
if (assembly.getSuppressedDuplicates() > 0) {
sb.append(assembly.getSuppressedDuplicates()).append(" duplicate");
if (assembly.getSuppressedDuplicates() != 1) sb.append("s");
sb.append(" suppressed");
}
if (assembly.getFilteredAcknowledged() > 0) {
if (assembly.getSuppressedDuplicates() > 0) sb.append(", ");
sb.append(assembly.getFilteredAcknowledged()).append(" acknowledged item");
if (assembly.getFilteredAcknowledged() != 1) sb.append("s");
sb.append(" filtered");
}
sb.append("_\n");
}
}
private void formatInsight(StringBuilder sb, DigestInsight insight) {
String emoji = getInsightEmoji(insight);
sb.append(emoji).append(" *").append(insight.getHeadline()).append("*\n");
if (insight.getDescription() != null && !insight.getDescription().isBlank()) {
sb.append(" ").append(insight.getDescription()).append("\n");
}
if (insight.isActionable() && insight.getSuggestedAction() != null) {
sb.append(" → _").append(insight.getSuggestedAction()).append("_\n");
}
// Add signature for dedup tracking
if (insight.getSignatureKey() != null) {
sb.append(" [sig:").append(insight.getSignatureKey()).append("]\n");
}
}
private String getInsightEmoji(DigestInsight insight) {
return switch (insight.getCategory()) {
case QUERY_PERFORMANCE -> "🐢";
case INDEX_RECOMMENDATIONS -> "📈";
case SCHEMA_CHANGES -> "🔧";
case GROWTH_ANOMALIES -> "📊";
case LOCK_CONCURRENCY -> "🔒";
case CONFIG_TUNING -> "⚙️";
case BRAIN_INTELLIGENCE -> "🧠";
case SYSTEM_ALERTS -> "🔔";
case COST_CAPACITY -> "💰";
case DOCUMENTATION_GAPS -> "📝";
};
}
/**
* Get the timestamp of the last digest sent for a connection.
* Used as the window start for personalized digest assembly.
*/
private LocalDateTime getLastDigestTime(String connectionId) {
return digestLogRepository
.findTopByConnectionIdOrderBySentAtDesc(connectionId)
.map(SlackDigestLog::getSentAt)
.orElse(LocalDateTime.now().minusHours(24));
}
private static final Duration DIGEST_TICK_LOOKBACK = Duration.ofMinutes(1);
/**
* Minute-tick entry point used by {@code SlackDailyDigestTaskConfig}.
*
* <p>When no enabled preferences exist, runs the legacy channel broadcast
* only if the global {@code slack.daily-digest.cron} is due (UTC).
* When preferences exist, delivers only to preferences whose cron matches
* in that user's timezone and that have not already been logged for this
* fire window.
*/
public void processDigestTick() {
processDigestTick(Instant.now());
}
/**
* Testable overload of {@link #processDigestTick()}.
*/
public void processDigestTick(Instant now) {
if (now == null) {
now = Instant.now();
}
long enabledCount = userDigestPreferenceRepository.countByEnabledTrue();
if (enabledCount == 0) {
if (DigestCronMatcher.isDue(globalDigestCron, ZoneId.of("UTC"), now, DIGEST_TICK_LOOKBACK)) {
log.info("Digest tick: no enabled preferences; running legacy broadcast (global cron due)");
runLegacyBroadcastForAllConnections();
} else {
log.debug("Digest tick: no enabled preferences; global cron not due");
}
return;
}
String globalCron = (globalDigestCron == null || globalDigestCron.isBlank())
? "0 0 9 * * *"
: globalDigestCron;
List<UserDigestPreference> enabledPrefs = userDigestPreferenceRepository
.findByEnabledTrueAndDeliveryMethod(DigestDeliveryMethod.SLACK_DM);
int dueCount = 0;
int delivered = 0;
for (UserDigestPreference pref : enabledPrefs) {
String cron = pref.getEffectiveCronExpression(globalCron);
ZoneId zone = DigestCronMatcher.resolveZone(pref.getTimezone());
var window = DigestCronMatcher.dueWindowStart(cron, zone, now, DIGEST_TICK_LOOKBACK);
if (window.isEmpty()) {
continue;
}
dueCount++;
Instant fireInstant = window.get();
List<String> connectionIds = resolvePreferenceConnections(pref);
for (String connectionId : connectionIds) {
if (alreadyDeliveredForWindow(pref, connectionId, fireInstant)) {
log.debug("Skipping digest for {} / {} — already delivered this window",
pref.getUsername(), connectionId);
continue;
}
try {
boolean slackEnabled = isSlackDeliveryEnabled();
String connName = connectionName(connectionId);
LocalDateTime since = getLastDigestTime(connectionId);
PersonalizedDigestResult result = sendPersonalizedDigestToUser(
connectionId, connName, pref, since, slackEnabled);
if (result.generated() || result.sent()) {
delivered++;
}
} catch (Exception e) {
log.error("Failed due-preference digest for user {} connection {}: {}",
pref.getUsername(), connectionId, e.getMessage(), e);
}
}
}
log.info("Digest tick (per-user): {} due preference(s), {} delivery attempt(s)",
dueCount, delivered);
// Connections with no recipients still get legacy broadcast when the global cron fires.
if (DigestCronMatcher.isDue(globalCron, ZoneId.of("UTC"), now, DIGEST_TICK_LOOKBACK)) {
for (String connectionId : digestConnectionIds()) {
if (userDigestPreferenceRepository.findEnabledForConnection(connectionId).isEmpty()) {
try {
sendLegacyDigest(connectionId);
} catch (Exception e) {
log.error("Legacy fallback digest failed for {}: {}", connectionId, e.getMessage(), e);
}
}
}
}
}
private void runLegacyBroadcastForAllConnections() {
List<String> connectionIds = digestConnectionIds();
if (connectionIds.isEmpty()) {
log.info("No connections available — skipping legacy digest");
return;
}
for (String connectionId : connectionIds) {
try {
sendLegacyDigest(connectionId);
} catch (Exception e) {
log.error("Failed legacy digest for connection {}: {}", connectionId, e.getMessage(), e);
}
}
}
private List<String> resolvePreferenceConnections(UserDigestPreference pref) {
if (pref.getConnectionId() != null && !pref.getConnectionId().isBlank()) {
return List.of(pref.getConnectionId());
}
return digestConnectionIds();
}
/**
* True when SlackDigestLog already has a row for this preference/connection
* at or after the cron fire time (idempotent minute-tick).
*/
private boolean alreadyDeliveredForWindow(
UserDigestPreference pref,
String connectionId,
Instant fireInstant) {
// sentAt is LocalDateTime.now() (JVM default zone) — compare in that zone.
LocalDateTime since = LocalDateTime.ofInstant(fireInstant, ZoneId.systemDefault());
if (pref.getId() != null) {
if (digestLogRepository.existsByPreferenceIdAndConnectionIdAndSentAtGreaterThanEqual(
pref.getId(), connectionId, since)) {
return true;
}
}
if (pref.getUsername() != null && !pref.getUsername().isBlank()) {
return digestLogRepository.existsByConnectionIdAndRecipientUsernameAndSentAtGreaterThanEqual(
connectionId, pref.getUsername(), since);
}
return false;
}
/**
* Entry point for the hybrid digest run: per-user when preferences exist,
* legacy broadcast otherwise.
*
* <p>Used by manual/admin triggers. The scheduled tick uses
* {@link #processDigestTick()} so per-user crons are honored.
*/
public void sendDailyDigestHybrid() {
List<String> connectionIds = digestConnectionIds();
if (connectionIds.isEmpty()) {
log.info("No connections available — skipping daily digest generation");
return;
}