Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
48 commits
Select commit Hold shift + click to select a range
8172c72
chore: WebSocket 의존성 및 CORS 설정 추가
JangInho May 17, 2026
57c0d97
feat: TrackingPoint 도메인 모델 추가
JangInho May 17, 2026
c0509d0
feat: 트래킹 GPS 실시간 수집 및 통계 구현
JangInho May 17, 2026
41ca956
fix: flush 트랜잭션 적용 위해 별도 서비스로 분리
JangInho May 19, 2026
7c7d219
feat: AGENTS.md 추가
pooreumjung May 21, 2026
2f198bd
Fix user withdraw persistence
pooreumjung May 21, 2026
a1cc3c2
탈퇴 후 재가입 시 새 사용자 생성
pooreumjung May 21, 2026
5199719
탈퇴 시 소셜 식별자 익명화
pooreumjung May 21, 2026
7f3e238
탈퇴 유저 로그인 방어 로직 추가
pooreumjung May 21, 2026
8397428
Merge pull request #64 from SEMOSAN/refactor/#62-user-withdraw
pooreumjung May 21, 2026
594b6c0
chore: update image to 839742812054d5de0ba8c02ecff975a9984b3db7 [skip…
github-actions[bot] May 21, 2026
f34a536
fix: 관악산 중복 row 정리 (V3 적용 전 hand-seed row 삭제)
JangInho May 21, 2026
f544939
Merge pull request #67 from SEMOSAN/fix/#65-legacy-gwanaksan
JangInho May 21, 2026
7ab3ca3
chore: update image to f5449397cc88ba378597b4bebd002889cddb31db [skip…
github-actions[bot] May 21, 2026
a189be8
chore: postgres 이미지를 postgis/postgis:16-alpine으로 변경
pooreumjung May 21, 2026
8f6a244
chore: update image to a189be835894efa3fd569a3f46e56ca3df6cd64d [skip…
github-actions[bot] May 21, 2026
349e814
카카오 로그인 액세스 토큰 방식 적용
pooreumjung May 21, 2026
445e75d
fix: postgres image 수정
pooreumjung May 21, 2026
f506a72
chore: update image to 445e75d50c4682fc5a980527201c23632994b26f [skip…
github-actions[bot] May 21, 2026
0b145de
Merge pull request #68 from SEMOSAN/feat/#18-tracking-session
JangInho May 21, 2026
e5beefb
chore: update image to 0b145deb46634b781209a4720d567c2e9fa67824 [skip…
github-actions[bot] May 21, 2026
569514d
fix: postgres 이미지를 ARM64 지원 버전으로 변경 (16-3.5)
pooreumjung May 21, 2026
ace41a6
chore: update image to 569514d987bca9b8539917eec3dd681569c51c46 [skip…
github-actions[bot] May 21, 2026
ff6beb3
fix: ARM64 지원 postgis 이미지로 변경 (imresamu/postgis:16-3.4)
pooreumjung May 21, 2026
324c6ea
chore: update image to ff6beb327877599b496fe3a74f16040138f31863 [skip…
github-actions[bot] May 21, 2026
457fafb
fix: V6 마이그레이션 외래키 제약 위반 수정 (amenities 먼저 삭제)
pooreumjung May 21, 2026
0cb828a
chore: update image to 457fafb40cc1ddc4ad26d766b77b77df42fe92fe [skip…
github-actions[bot] May 21, 2026
c274afe
fix: V6 마이그레이션 외래키 제약 위반 수정 (mountains 참조 테이블 cascade 삭제)
pooreumjung May 21, 2026
cb43205
chore: update image to c274afe2fc9f6457f4e20fa2b5dd50ec4191a4bf [skip…
github-actions[bot] May 21, 2026
5819cdb
fix: V6 마이그레이션 restaurants 외래키 제약 위반 추가 수정
pooreumjung May 21, 2026
2598403
chore: update image to 5819cdbe39cd86400117feef8ebb9cfd6541a088 [skip…
github-actions[bot] May 21, 2026
1fb1de7
fix: HandlerMethodValidationException 핸들러 중복 충돌 수정 (@Override로 변경)
pooreumjung May 21, 2026
cbbda93
chore: update image to 1fb1de7dc7c67566c54a61674c69c1ff7d7685fa [skip…
github-actions[bot] May 21, 2026
d27b7b8
Merge pull request #69 from SEMOSAN/refactor/#66-kakao-login
pooreumjung May 21, 2026
73bc7e0
chore: update image to d27b7b86343afd7846e858d29cd5d0bdc8c41f25 [skip…
github-actions[bot] May 21, 2026
c0d6d99
feat: 기본 이름 생성 리소스 추가
pooreumjung May 21, 2026
8d76665
feat: 신규 유저 기본 이름 자동 생성
pooreumjung May 21, 2026
4645a2f
feat: 유저 프로필 응답에 이름 추가
pooreumjung May 21, 2026
99bf9c9
fix: 기본 표시값을 닉네임에 적용
pooreumjung May 21, 2026
a88e749
refactor: 기본 닉네임 생성 검증 개선
pooreumjung May 21, 2026
2408eb0
Merge pull request #72 from SEMOSAN/feat/#71-generate-user-basic-name
pooreumjung May 21, 2026
9a7bb76
chore: update image to 2408eb0c0f5fc413f174a8745cfd08f32d7a88ce [skip…
github-actions[bot] May 21, 2026
e3897f6
Merge remote-tracking branch 'origin/develop' into feat/#19-tracking-gps
JangInho May 21, 2026
ee96567
Merge remote-tracking branch 'origin/develop' into feat/#19-tracking-gps
JangInho May 21, 2026
7054f1e
refactor: GPS 점 flush 실패 시 유실 방지 및 recordedAt 검증 추가
JangInho May 21, 2026
0b2b4cc
refactor: 트래킹 세션 종료 시 GPS 버퍼 정리 (메모리 누수 방지)
JangInho May 21, 2026
ffa4cfc
refactor: 스케줄러 stale 판정 기준을 Redis 활성 마커로 변경
JangInho May 21, 2026
c227637
chore: tracking_points 테이블 마이그레이션 추가
JangInho May 21, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 42 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
# AGENTS.md

## Communication
- Always reply in Korean unless the user explicitly asks for another language.
- Keep answers direct and practical. Avoid repeating the same explanation after the user has already acknowledged it.
- When the user asks "확인해줘", "봐줘", "어떻게 돼", "뭐가 문제야", or similar, treat it as analysis-only unless they explicitly ask for code changes.
- Do not edit files unless the user clearly says something like "수정해줘", "고쳐줘", "반영해줘", "만들어줘", "추가해줘", or "삭제해줘".
- If a fix is likely, explain the cause, affected files, and proposed patch first. Wait for explicit approval before applying it.
- If you accidentally changed code without approval, say exactly which files changed and offer to revert only your own changes.

## User Preferences
- The user wants investigation before implementation.
- The user dislikes unrequested code edits.
- The user prefers concrete answers based on the current codebase, not generic guesses.
- When diagnosing backend/frontend integration issues, clearly separate:
- what the backend currently expects
- what the frontend must send
- what environment/config values must match
- what is only an assumption
- If sensitive values such as secrets, tokens, DB passwords, or JWT secrets appear in screenshots or logs, warn briefly and recommend rotation if they may have been shared.

## Project
- This is a Spring Boot backend project.
- Use Java 21 and Gradle.
- Follow the existing package structure, naming, and style.
- Prefer existing services, repositories, DTOs, and response conventions over introducing new abstractions.

## Commands
- Use `rg` first for searching.
- Run focused tests before broad tests when checking a narrow change.
- Common commands:
- `./gradlew test --tests <fully.qualified.TestClass>`
- `./gradlew test`
- `./gradlew build`
- If full tests fail because local infrastructure such as PostgreSQL or Redis is unavailable, report that clearly and distinguish it from failures caused by the code change.

## Editing Rules
- Never revert or overwrite user changes unless explicitly requested.
- Keep changes narrowly scoped to the requested task.
- Do not commit secrets, environment values, generated local files, or unrelated formatting churn.
- Before editing, state what files will be touched and why.
- After editing, summarize the exact files changed and verification performed.
4 changes: 4 additions & 0 deletions build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -53,9 +53,13 @@ dependencies {
// PostGIS (Hibernate Spatial + JTS)
implementation 'org.hibernate.orm:hibernate-spatial'

// WebSocket (트래킹 GPS 실시간 수집 — STOMP)
implementation 'org.springframework.boot:spring-boot-starter-websocket'

// Flyway
implementation 'org.flywaydb:flyway-core'
implementation 'org.flywaydb:flyway-database-postgresql'

}

jar { enabled = false }
Expand Down
2 changes: 1 addition & 1 deletion k8s/kustomization.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -14,4 +14,4 @@ resources:
- minio/service.yaml
images:
- name: ghcr.io/semosan/semosan_be
newTag: fb0eb2823ad7e4ee955bd3941443be46e88a22cf
newTag: 2408eb0c0f5fc413f174a8745cfd08f32d7a88ce
4 changes: 2 additions & 2 deletions k8s/postgres/deployment.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ spec:
spec:
containers:
- name: postgres
image: postgres:16-alpine
image: imresamu/postgis:16-3.4
ports:
- containerPort: 5432
envFrom:
Expand All @@ -28,4 +28,4 @@ spec:
volumes:
- name: postgres-data
persistentVolumeClaim:
claimName: postgres-pvc
claimName: postgres-pvc
20 changes: 20 additions & 0 deletions src/main/java/com/semosan/api/common/config/SecurityConfig.java
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,15 @@ public class SecurityConfig {
"/api/auth/token/reissue"
};

/**
* WebSocket(STOMP) 엔드포인트.
* 핸드셰이크는 HTTP JWT 필터로 인증하지 않고 통과시키고,
* STOMP CONNECT 프레임에서 StompAuthChannelInterceptor 가 JWT 를 검증한다.
*/
public static final String[] WEBSOCKET_URIS = {
"/ws/tracking/**"
};

@Bean
public SecurityFilterChain securityFilterChain(HttpSecurity http) throws Exception {
http
Expand All @@ -63,6 +72,7 @@ public SecurityFilterChain securityFilterChain(HttpSecurity http) throws Excepti
.requestMatchers(SWAGGER_URIS).permitAll()
.requestMatchers(OAUTH_URIS).permitAll()
.requestMatchers(AUTH_URIS).permitAll()
.requestMatchers(WEBSOCKET_URIS).permitAll()
.anyRequest().authenticated()
)
.addFilterBefore(jwtFilter, UsernamePasswordAuthenticationFilter.class);
Expand All @@ -78,7 +88,17 @@ public CorsConfigurationSource corsConfigurationSource() {
config.setAllowedHeaders(List.of("*"));
config.setAllowCredentials(true);

// WebSocket(STOMP) 핸드셰이크는 다양한 origin(모바일/로컬 테스트 페이지 등)에서 들어옴.
// /ws/** 만 별도 정책으로 풀어준다.
// TODO: production 에서는 모바일 앱 origin 만 명시적으로 허용하도록 좁힐 것.
CorsConfiguration wsConfig = new CorsConfiguration();
wsConfig.setAllowedOriginPatterns(List.of("*"));
wsConfig.setAllowedMethods(List.of("GET"));
wsConfig.setAllowedHeaders(List.of("*"));
wsConfig.setAllowCredentials(true);

Comment on lines +91 to +99

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟠 Major | ⚡ Quick win

Restrict WebSocket origins now (don’t ship wildcard + credentials).

Line 95 currently allows any origin for /ws/** while Line 98 enables credentials. This leaves the handshake policy overly permissive for a privileged channel. Replace * with an explicit allowlist and keep it environment-driven. Also align WebSocketConfig Line 29 to the same allowlist.

🔒 Suggested change
-        wsConfig.setAllowedOriginPatterns(List.of("*"));
+        wsConfig.setAllowedOrigins(List.of(
+                "http://localhost:3000",
+                "https://lgenius.site"
+        ));

Also applies to: 101-101

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@src/main/java/com/semosan/api/common/config/SecurityConfig.java` around lines
91 - 99, The WebSocket CORS config is too permissive: update the
CorsConfiguration instance named wsConfig to use an environment-driven explicit
allowlist instead of setAllowedOriginPatterns(List.of("*")), keep
setAllowCredentials(true) only if the allowlist is not wildcard, and ensure
setAllowedOriginPatterns (or setAllowedOrigins) is populated from a configurable
property (e.g., a comma-separated env/property) so production can restrict to
mobile app origins; also mirror the same change in WebSocketConfig (Line 29) so
both places read the same allowlist property.

UrlBasedCorsConfigurationSource source = new UrlBasedCorsConfigurationSource();
source.registerCorsConfiguration("/ws/**", wsConfig);
source.registerCorsConfiguration("/**", config);
return source;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,13 +7,6 @@
@Configuration
public class WebClientConfig {

@Bean
public WebClient kakaoAuthWebClient() {
return WebClient.builder()
.baseUrl("https://kauth.kakao.com")
.build();
}

@Bean
public WebClient kakaoApiWebClient() {
return WebClient.builder()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -55,16 +55,20 @@ public ResponseEntity<ApiResponse<Void>> handleConstraintViolation(
}

// Spring 6.1+ HandlerMethodValidationException (컨트롤러 메서드 파라미터 검증 실패 표준 경로)
@ExceptionHandler(HandlerMethodValidationException.class)
public ResponseEntity<ApiResponse<Void>> handleHandlerMethodValidation(
HandlerMethodValidationException e
@Override
protected ResponseEntity<Object> handleHandlerMethodValidationException(
HandlerMethodValidationException ex,
HttpHeaders headers,
HttpStatusCode status,
WebRequest request
) {
String message = e.getAllErrors().stream()
String message = ex.getAllErrors().stream()
.findFirst()
.map(error -> error.getDefaultMessage() != null ? error.getDefaultMessage() : ErrorStatus.BAD_REQUEST.getMessage())
.orElse(ErrorStatus.BAD_REQUEST.getMessage());
log.warn("[*] HandlerMethodValidationException : {}", message);
return ApiResponse.error(ErrorStatus.BAD_REQUEST, message);
ApiResponse<Void> body = createApiResponse(ErrorStatus.BAD_REQUEST, message);
return handleExceptionInternal(ex, body, headers, status, request);
}

// null 참조로 발생한 서버 오류를 500 에러로 응답
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,15 +3,11 @@
import com.semosan.api.common.exception.GeneralException;
import com.semosan.api.domain.oauth.properties.KakaoProperties;
import com.semosan.api.common.status.ErrorStatus;
import com.semosan.api.domain.oauth.dto.KakaoTokenResponse;
import com.semosan.api.domain.oauth.dto.KakaoUserInfoResponse;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.http.HttpStatusCode;
import org.springframework.stereotype.Component;
import org.springframework.util.LinkedMultiValueMap;
import org.springframework.util.MultiValueMap;
import org.springframework.util.StringUtils;
import org.springframework.web.reactive.function.BodyInserters;
import org.springframework.web.reactive.function.client.WebClient;
import org.springframework.web.reactive.function.client.WebClientResponseException;
Expand All @@ -22,44 +18,9 @@
@RequiredArgsConstructor
public class OAuthKakaoClient {

private final WebClient kakaoAuthWebClient;
private final WebClient kakaoApiWebClient;
private final KakaoProperties kakaoProperties;

// 인가 코드 -> 카카오 액세스 토큰
public KakaoTokenResponse getKakaoToken(String code) {
MultiValueMap<String, String> form = new LinkedMultiValueMap<>();
form.add("grant_type", "authorization_code");
form.add("client_id", kakaoProperties.clientId());
form.add("redirect_uri", kakaoProperties.redirectUri());
form.add("code", code);
if (StringUtils.hasText(kakaoProperties.clientSecret())) {
form.add("client_secret", kakaoProperties.clientSecret());
}

try {
KakaoTokenResponse response = kakaoAuthWebClient.post()
.uri("/oauth/token")
.body(BodyInserters.fromFormData(form))
.retrieve()
.onStatus(
HttpStatusCode::isError,
r -> r.bodyToMono(String.class)
.flatMap(body -> handleError(r.statusCode(), body, "토큰 발급"))
)
.bodyToMono(KakaoTokenResponse.class)
.block();

if (response == null || !StringUtils.hasText(response.accessToken())) {
throw new GeneralException(ErrorStatus.KAKAO_TOKEN_REQUEST_FAILED);
}
return response;
} catch (WebClientResponseException e) {
log.error("[*] 카카오 토큰 발급 실패 status={}, body={}", e.getStatusCode(), e.getResponseBodyAsString());
throw new GeneralException(ErrorStatus.KAKAO_TOKEN_REQUEST_FAILED);
}
}

// 카카오 액세스 토큰 -> 사용자 정보
public KakaoUserInfoResponse getKakaoUserInfo(String kakaoAccessToken) {
try {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ public interface OAuthControllerDocs {

@Operation(
summary = "카카오 소셜 로그인",
description = "프론트엔드에서 전달받은 카카오 인가 코드(code)로 로그인 또는 회원가입을 처리하고 서비스 JWT를 발급합니다."
description = "프론트엔드에서 전달받은 카카오 액세스 토큰(accessToken)으로 로그인 또는 회원가입을 처리하고 서비스 JWT를 발급합니다."
)
@ApiResponses({
@io.swagger.v3.oas.annotations.responses.ApiResponse(
Expand All @@ -28,7 +28,7 @@ public interface OAuthControllerDocs {
),
@io.swagger.v3.oas.annotations.responses.ApiResponse(
responseCode = "400",
description = "잘못된 요청 (인가 코드 또는 디바이스 타입 누락)",
description = "잘못된 요청 (카카오 액세스 토큰 또는 디바이스 타입 누락)",
content = @Content(schema = @Schema(implementation = ApiResponse.class))
),
@io.swagger.v3.oas.annotations.responses.ApiResponse(
Expand Down

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,8 @@

public record OAuthKakaoLoginRequest(

@NotBlank(message = "인가 코드는 필수입니다.")
String code,
@NotBlank(message = "카카오 액세스 토큰은 필수입니다.")
String accessToken,

@NotNull(message = "디바이스 타입은 필수입니다.")
DeviceType deviceType
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,6 @@
import com.semosan.api.common.jwt.TokenIssuance;
import com.semosan.api.domain.oauth.client.OAuthAppleClient;
import com.semosan.api.domain.oauth.client.OAuthKakaoClient;
import com.semosan.api.domain.oauth.dto.KakaoTokenResponse;
import com.semosan.api.domain.oauth.dto.KakaoUserInfoResponse;
import com.semosan.api.domain.oauth.dto.request.OAuthAppleLoginRequest;
import com.semosan.api.domain.oauth.dto.request.OAuthKakaoLoginRequest;
Expand All @@ -29,8 +28,7 @@ public class OAuthService {

@Transactional
public OAuthLoginResponse kakaoLogin(OAuthKakaoLoginRequest request) {
KakaoTokenResponse kakaoToken = oAuthKakaoClient.getKakaoToken(request.code());
KakaoUserInfoResponse userInfo = oAuthKakaoClient.getKakaoUserInfo(kakaoToken.accessToken());
KakaoUserInfoResponse userInfo = oAuthKakaoClient.getKakaoUserInfo(request.accessToken());

// DTO 파싱은 oauth 레이어에서 처리 후 순수 값만 UserService로 전달
KakaoUserInfoResponse.KakaoAccount account = userInfo.kakaoAccount();
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,75 @@
package com.semosan.api.domain.tracking.config;

import com.semosan.api.common.config.TrackingProperties;
import com.semosan.api.domain.tracking.service.TrackingStreamConsumer;
import jakarta.annotation.PreDestroy;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.redis.connection.RedisConnectionFactory;
import org.springframework.data.redis.connection.stream.Consumer;
import org.springframework.data.redis.connection.stream.MapRecord;
import org.springframework.data.redis.connection.stream.ReadOffset;
import org.springframework.data.redis.connection.stream.StreamOffset;
import org.springframework.data.redis.stream.StreamMessageListenerContainer;

import java.net.InetAddress;
import java.net.UnknownHostException;
import java.time.Duration;
import java.util.UUID;

/**
* Redis Stream(tracking:gps) 의 GPS 점을 소비하는 컨테이너를 부트업한다.
* - Consumer group 자체는 RedisStreamConfig 에서 사전 생성됨.
* - 본 클래스는 그룹의 이 인스턴스를 위한 consumer 등록 + listener container 시작.
* - consumer 이름은 호스트명 + UUID prefix 로 인스턴스별 유일성 확보 (다중 인스턴스 대비).
*/
@Slf4j
@Configuration
@RequiredArgsConstructor
public class TrackingStreamListenerConfig implements InitializingBean {

private final RedisConnectionFactory redisConnectionFactory;
private final TrackingProperties trackingProperties;
private final TrackingStreamConsumer trackingStreamConsumer;

private StreamMessageListenerContainer<String, MapRecord<String, String, String>> container;

@Override
public void afterPropertiesSet() {
var options = StreamMessageListenerContainer
.StreamMessageListenerContainerOptions.builder()
.pollTimeout(Duration.ofSeconds(1))
.build();

container = StreamMessageListenerContainer.create(redisConnectionFactory, options);

container.receiveAutoAck(
Consumer.from(trackingProperties.getConsumerGroup(), buildConsumerName()),
StreamOffset.create(trackingProperties.getStreamKey(), ReadOffset.lastConsumed()),
trackingStreamConsumer
);
Comment on lines +48 to +52

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟠 Major | 🏗️ Heavy lift

🧩 Analysis chain

🌐 Web query:

In Spring Data Redis StreamMessageListenerContainer, does receiveAutoAck acknowledge entries before listener processing succeeds, and what pattern is recommended for at-least-once processing with consumer groups?

💡 Result:

Yes. In StreamMessageListenerContainer, using receiveAutoAck (auto-acknowledge) acknowledges messages when they are received, not after your listener code completes successfully. This is explicitly stated as “Every message is acknowledged when received” for receiveAutoAck, whereas the non-auto-ack mode requires acknowledging via StreamOperations.acknowledge after processing completes successfully [1][2][3]. Recommended pattern for at-least-once with consumer groups: 1) Use external acknowledgement: call receive(...) (autoAck=false) and only acknowledge with StreamOperations.acknowledge after your listener processing succeeds [1][4][3]. This keeps failed/unacknowledged deliveries in the consumer group’s Pending Entries List (PEL) [3][5]. 2) Handle PENDING recovery explicitly: start (or periodically run) a “drain PEL” phase for your consumer group (e.g., read pending/previously delivered entries), and only when the PEL is empty switch to reading new messages (ID “>”). Redis documents that re-delivered pending entries remain pending until explicitly acknowledged via XACK, and the typical recovery approach is to first process the pending backlog (using an initial ID such as 0), then proceed to “>” for new messages [6]. In code terms, the container doesn’t provide a built-in “always re-process PENDING until it succeeds before polling for new” switch; you implement the two-phase (PEL-drain then new) flow by choosing the appropriate StreamRead offsets/requests (PEL history first, then “>”), and by acknowledging only after successful processing [1][3][6].

Citations:


Use manual acknowledgement instead of auto-ack to ensure at-least-once delivery for GPS points.

receiveAutoAck() acknowledges entries when received, before your listener processes them. If downstream processing fails, acknowledged entries are permanently lost and cannot be retried from the consumer group's pending backlog.

Switch to receive() (non-auto-ack) and call StreamOperations.acknowledge() only after processing succeeds. This keeps failed entries in the Pending Entries List for recovery and retries.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In
`@src/main/java/com/semosan/api/domain/tracking/config/TrackingStreamListenerConfig.java`
around lines 48 - 52, The code currently uses container.receiveAutoAck(...)
which acknowledges entries immediately; change to container.receive(...)
(non-auto-ack) using the same
Consumer.from(trackingProperties.getConsumerGroup(), buildConsumerName()),
StreamOffset.create(trackingProperties.getStreamKey(),
ReadOffset.lastConsumed()), and trackingStreamConsumer so records are delivered
but not auto-acked. Then update your trackingStreamConsumer processing logic to
call the Redis stream acknowledge API (e.g., StreamOperations.acknowledge(...)
or the container/RedisTemplate opsForStream().acknowledge) after successful
processing using trackingProperties.getConsumerGroup() and the record id; do not
acknowledge on failure so entries remain in the Pending Entries List for
retries.

container.start();
log.info("Started Redis Stream listener: stream={} group={}",
trackingProperties.getStreamKey(),
trackingProperties.getConsumerGroup());
}

@PreDestroy
public void stop() {
if (container != null) {
container.stop();
}
}

private static String buildConsumerName() {
String host;
try {
host = InetAddress.getLocalHost().getHostName();
} catch (UnknownHostException e) {
host = "unknown";
}
return host + "-" + UUID.randomUUID().toString().substring(0, 8);
}
}
Loading