Skip to content

Commit 731da57

Browse files
Refactored
1 parent 2466827 commit 731da57

3 files changed

Lines changed: 122 additions & 156 deletions

File tree

src/main/java/com/byteentropy/idempotency_core/aspect/IdempotencyAspect.java

Lines changed: 18 additions & 97 deletions
Original file line numberDiff line numberDiff line change
@@ -3,11 +3,7 @@
33
import com.byteentropy.idempotency_core.annotation.Idempotent;
44
import com.byteentropy.idempotency_core.model.IdempotencyRecord;
55
import com.byteentropy.idempotency_core.model.IdempotencyStatus;
6-
import com.byteentropy.idempotency_core.storage.IdempotencyStore;
7-
import com.fasterxml.jackson.core.JsonProcessingException;
8-
import com.fasterxml.jackson.databind.DeserializationFeature;
9-
import com.fasterxml.jackson.databind.ObjectMapper;
10-
import com.fasterxml.jackson.databind.SerializationFeature;
6+
import com.byteentropy.idempotency_core.service.IdempotencyService;
117
import org.aspectj.lang.ProceedingJoinPoint;
128
import org.aspectj.lang.annotation.Around;
139
import org.aspectj.lang.annotation.Aspect;
@@ -23,34 +19,21 @@
2319
import org.springframework.stereotype.Component;
2420
import org.springframework.util.StringUtils;
2521

26-
import java.lang.reflect.Method;
27-
import java.nio.charset.StandardCharsets;
28-
import java.security.MessageDigest;
29-
import java.security.NoSuchAlgorithmException;
30-
import java.util.HexFormat;
3122
import java.util.Objects;
3223

3324
@Aspect
3425
@Component
3526
public class IdempotencyAspect implements Ordered {
36-
3727
private static final Logger log = LoggerFactory.getLogger(IdempotencyAspect.class);
38-
private final IdempotencyStore store;
28+
29+
private final IdempotencyService idempotencyService;
3930
private final ExpressionParser parser = new SpelExpressionParser();
40-
private final ObjectMapper hashMapper;
4131

4232
@Value("${idempotency.default-ttl:3600}")
4333
private long globalDefaultTtl;
4434

45-
@Value("${idempotency.processing-timeout-ms:300000}")
46-
private long processingTimeoutMs;
47-
48-
public IdempotencyAspect(IdempotencyStore store) {
49-
this.store = store;
50-
// Pre-configure a safe mapper for hashing to avoid GC pressure from frequent .copy() calls
51-
this.hashMapper = store.getObjectMapper().copy()
52-
.disable(SerializationFeature.FAIL_ON_EMPTY_BEANS)
53-
.disable(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES);
35+
public IdempotencyAspect(IdempotencyService idempotencyService) {
36+
this.idempotencyService = idempotencyService;
5437
}
5538

5639
@Override
@@ -61,77 +44,43 @@ public int getOrder() {
6144
@Around("@annotation(com.byteentropy.idempotency_core.annotation.Idempotent)")
6245
public Object handle(ProceedingJoinPoint joinPoint) throws Throwable {
6346
MethodSignature signature = (MethodSignature) joinPoint.getSignature();
64-
Method method = signature.getMethod();
65-
Idempotent idempotent = method.getAnnotation(Idempotent.class);
47+
Idempotent idempotent = signature.getMethod().getAnnotation(Idempotent.class);
6648

67-
return execute(joinPoint, idempotent);
68-
}
69-
70-
private Object execute(ProceedingJoinPoint joinPoint, Idempotent idempotent) throws Throwable {
7149
String namespace = idempotent.namespace();
7250
String key = resolveSpel(joinPoint, idempotent.key());
7351
long finalTtl = (idempotent.ttl() > 0) ? idempotent.ttl() : globalDefaultTtl;
7452

75-
// Validation
76-
if (!StringUtils.hasText(namespace)) {
77-
throw new IllegalArgumentException("Idempotency namespace is required.");
78-
}
79-
if (!StringUtils.hasText(key)) {
80-
throw new IllegalArgumentException("Idempotency key evaluated to empty/null.");
53+
if (!StringUtils.hasText(namespace) || !StringUtils.hasText(key)) {
54+
throw new IllegalArgumentException("Idempotency namespace and key are required.");
8155
}
8256

83-
String currentRequestHash = generateRequestHash(joinPoint.getArgs());
57+
String currentHash = idempotencyService.generateHash(joinPoint.getArgs());
58+
IdempotencyRecord existing = idempotencyService.attemptReservation(namespace, key, currentHash, finalTtl);
8459

85-
IdempotencyRecord initial = IdempotencyRecord.builder()
86-
.status(IdempotencyStatus.PROCESSING)
87-
.requestHash(currentRequestHash)
88-
.timestamp(System.currentTimeMillis())
89-
.build();
90-
91-
// Atomic check-and-reserve
92-
Object resultFromStore = store.executeLua(namespace, key, initial, finalTtl);
93-
94-
if (resultFromStore != null) {
95-
IdempotencyRecord existing = (IdempotencyRecord) resultFromStore;
96-
97-
// 1. Handle concurrent processing
60+
if (existing != null) {
9861
if (existing.getStatus() == IdempotencyStatus.PROCESSING) {
9962
long elapsed = System.currentTimeMillis() - existing.getTimestamp();
100-
if (elapsed > processingTimeoutMs) {
101-
log.warn("Ghost lock detected for {}:{}. Clearing and retrying.", namespace, key);
102-
store.delete(namespace, key);
103-
return execute(joinPoint, idempotent);
63+
if (elapsed > idempotencyService.getProcessingTimeoutMs()) {
64+
log.warn("Ghost lock detected for {}:{}. Retrying.", namespace, key);
65+
idempotencyService.rollback(namespace, key);
66+
return handle(joinPoint);
10467
}
10568
throw new RuntimeException("Request is currently being processed.");
10669
}
10770

108-
// 2. Validate payload integrity (Same key, different data)
109-
if (!Objects.equals(existing.getRequestHash(), currentRequestHash)) {
71+
if (!Objects.equals(existing.getRequestHash(), currentHash)) {
11072
throw new IllegalStateException("Idempotency Conflict: Key exists with different payload.");
11173
}
11274

113-
// 3. Return cached response
114-
log.info("Returning cached response for {}:{}", namespace, key);
11575
return existing.getResponse();
11676
}
11777

11878
try {
119-
// Execute business logic
12079
Object response = joinPoint.proceed();
121-
122-
IdempotencyRecord completed = IdempotencyRecord.builder()
123-
.status(IdempotencyStatus.COMPLETED)
124-
.response(response)
125-
.requestHash(currentRequestHash)
126-
.timestamp(System.currentTimeMillis())
127-
.build();
128-
129-
store.save(namespace, key, completed, finalTtl);
80+
idempotencyService.commit(namespace, key, currentHash, response, finalTtl);
13081
return response;
13182
} catch (Throwable e) {
132-
// Clean up the lock on failure so the client can retry
133-
log.error("Execution failed for {}:{}. Clearing lock.", namespace, key);
134-
store.delete(namespace, key);
83+
idempotencyService.rollback(namespace, key);
13584
throw e;
13685
}
13786
}
@@ -150,32 +99,4 @@ private String resolveSpel(ProceedingJoinPoint joinPoint, String spel) {
15099
context.setVariable("methodName", signature.getMethod().getName());
151100
return parser.parseExpression(spel).getValue(context, String.class);
152101
}
153-
154-
private String generateRequestHash(Object[] args) {
155-
if (args == null || args.length == 0) return "no-args";
156-
157-
try {
158-
StringBuilder sb = new StringBuilder();
159-
for (Object arg : args) {
160-
if (arg == null) {
161-
sb.append("null");
162-
} else {
163-
sb.append(hashMapper.writeValueAsString(arg));
164-
}
165-
}
166-
167-
MessageDigest digest = MessageDigest.getInstance("SHA-256");
168-
byte[] encodedHash = digest.digest(sb.toString().getBytes(StandardCharsets.UTF_8));
169-
return HexFormat.of().formatHex(encodedHash);
170-
171-
} catch (NoSuchAlgorithmException e) {
172-
log.error("SHA-256 not available, falling back to identity hash.");
173-
return "fallback-id-" + Objects.hash(args);
174-
} catch (JsonProcessingException e) {
175-
log.warn("Serialization for hashing failed. Falling back to identity hash.");
176-
return "fallback-json-" + Objects.hash(args);
177-
} catch (Exception e) {
178-
return "hash-err-" + Objects.hash(args);
179-
}
180-
}
181102
}

src/main/java/com/byteentropy/idempotency_core/controller/IdempotencyController.java

Lines changed: 16 additions & 59 deletions
Original file line numberDiff line numberDiff line change
@@ -1,60 +1,39 @@
11
package com.byteentropy.idempotency_core.controller;
22

3-
import com.byteentropy.idempotency_core.model.*;
43
import com.byteentropy.idempotency_core.api.IdempotencyRequest;
54
import com.byteentropy.idempotency_core.api.IdempotencyResponse;
6-
import com.byteentropy.idempotency_core.storage.IdempotencyStore;
7-
import com.fasterxml.jackson.databind.DeserializationFeature;
8-
import com.fasterxml.jackson.databind.ObjectMapper;
9-
import com.fasterxml.jackson.databind.SerializationFeature;
5+
import com.byteentropy.idempotency_core.model.IdempotencyRecord;
6+
import com.byteentropy.idempotency_core.model.IdempotencyStatus;
7+
import com.byteentropy.idempotency_core.service.IdempotencyService;
108
import org.springframework.http.HttpStatus;
119
import org.springframework.http.ResponseEntity;
1210
import org.springframework.util.StringUtils;
1311
import org.springframework.web.bind.annotation.*;
1412

15-
import java.nio.charset.StandardCharsets;
16-
import java.security.MessageDigest;
17-
import java.util.HexFormat;
1813
import java.util.Objects;
1914

2015
@RestController
2116
@RequestMapping("/api/v1/idempotency")
2217
public class IdempotencyController {
2318

24-
private final IdempotencyStore store;
25-
private final ObjectMapper hashMapper;
19+
private final IdempotencyService idempotencyService;
2620

27-
public IdempotencyController(IdempotencyStore store, ObjectMapper objectMapper) {
28-
this.store = store;
29-
// Use a consistent, safe mapper for hashing
30-
this.hashMapper = objectMapper.copy()
31-
.disable(SerializationFeature.FAIL_ON_EMPTY_BEANS)
32-
.disable(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES);
21+
public IdempotencyController(IdempotencyService idempotencyService) {
22+
this.idempotencyService = idempotencyService;
3323
}
3424

3525
@PostMapping("/check")
36-
public ResponseEntity<IdempotencyResponse> check(@RequestBody IdempotencyRequest request) throws Exception {
37-
26+
public ResponseEntity<IdempotencyResponse> check(@RequestBody IdempotencyRequest request) {
3827
if (!StringUtils.hasText(request.namespace())) {
3928
return ResponseEntity.status(HttpStatus.BAD_REQUEST)
4029
.body(new IdempotencyResponse(request.key(), null, null, "Namespace is required", System.currentTimeMillis()));
4130
}
4231

43-
String ns = request.namespace();
44-
String currentHash = generateHash(request.payload());
45-
46-
IdempotencyRecord initial = IdempotencyRecord.builder()
47-
.status(IdempotencyStatus.PROCESSING)
48-
.requestHash(currentHash)
49-
.timestamp(System.currentTimeMillis())
50-
.build();
32+
String currentHash = idempotencyService.generateHash(request.payload());
33+
IdempotencyRecord existing = idempotencyService.attemptReservation(
34+
request.namespace(), request.key(), currentHash, request.ttl());
5135

52-
Object result = store.executeLua(ns, request.key(), initial, request.ttl());
53-
54-
if (result != null) {
55-
IdempotencyRecord existing = (IdempotencyRecord) result;
56-
57-
// Critical: Must compare SHA-256 hashes
36+
if (existing != null) {
5837
if (!Objects.equals(existing.getRequestHash(), currentHash)) {
5938
return ResponseEntity.status(HttpStatus.CONFLICT)
6039
.body(new IdempotencyResponse(request.key(), null, null, "Payload mismatch", System.currentTimeMillis()));
@@ -73,38 +52,16 @@ public ResponseEntity<IdempotencyResponse> check(@RequestBody IdempotencyRequest
7352
}
7453

7554
@PostMapping("/complete")
76-
public ResponseEntity<IdempotencyResponse> complete(@RequestBody CompletionRequest wrapper) throws Exception {
55+
public ResponseEntity<Void> complete(@RequestBody CompletionRequest wrapper) {
7756
IdempotencyRequest request = wrapper.request();
78-
79-
if (!StringUtils.hasText(request.namespace())) {
57+
if (request == null || !StringUtils.hasText(request.namespace())) {
8058
return ResponseEntity.status(HttpStatus.BAD_REQUEST).build();
8159
}
8260

83-
String ns = request.namespace();
84-
String hash = generateHash(request.payload());
85-
86-
IdempotencyRecord completed = IdempotencyRecord.builder()
87-
.status(IdempotencyStatus.COMPLETED)
88-
.response(wrapper.resultData())
89-
.requestHash(hash)
90-
.timestamp(System.currentTimeMillis())
91-
.build();
92-
93-
store.save(ns, request.key(), completed, request.ttl());
94-
return ResponseEntity.ok().build();
95-
}
96-
97-
/**
98-
* Internal helper to generate SHA-256 hash consistent with IdempotencyAspect
99-
*/
100-
private String generateHash(Object payload) throws Exception {
101-
if (payload == null) return "null-payload";
102-
103-
String json = hashMapper.writeValueAsString(payload);
104-
MessageDigest digest = MessageDigest.getInstance("SHA-256");
105-
byte[] encodedHash = digest.digest(json.getBytes(StandardCharsets.UTF_8));
61+
String hash = idempotencyService.generateHash(request.payload());
62+
idempotencyService.commit(request.namespace(), request.key(), hash, wrapper.resultData(), request.ttl());
10663

107-
return HexFormat.of().formatHex(encodedHash);
64+
return ResponseEntity.ok().build();
10865
}
10966

11067
public record CompletionRequest(IdempotencyRequest request, Object resultData) {}
Lines changed: 88 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,88 @@
1+
package com.byteentropy.idempotency_core.service;
2+
3+
import com.byteentropy.idempotency_core.model.IdempotencyRecord;
4+
import com.byteentropy.idempotency_core.model.IdempotencyStatus;
5+
import com.byteentropy.idempotency_core.storage.IdempotencyStore;
6+
import com.fasterxml.jackson.databind.DeserializationFeature;
7+
import com.fasterxml.jackson.databind.ObjectMapper;
8+
import com.fasterxml.jackson.databind.SerializationFeature;
9+
import org.slf4j.Logger;
10+
import org.slf4j.LoggerFactory;
11+
import org.springframework.beans.factory.annotation.Value;
12+
import org.springframework.stereotype.Service;
13+
14+
import java.nio.charset.StandardCharsets;
15+
import java.security.MessageDigest;
16+
import java.util.HexFormat;
17+
import java.util.Objects;
18+
19+
@Service
20+
public class IdempotencyService {
21+
private static final Logger log = LoggerFactory.getLogger(IdempotencyService.class);
22+
23+
private final IdempotencyStore store;
24+
private final ObjectMapper hashMapper;
25+
26+
@Value("${idempotency.processing-timeout-ms:300000}")
27+
private long processingTimeoutMs;
28+
29+
public IdempotencyService(IdempotencyStore store) {
30+
this.store = store;
31+
this.hashMapper = store.getObjectMapper().copy()
32+
.disable(SerializationFeature.FAIL_ON_EMPTY_BEANS)
33+
.disable(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES);
34+
}
35+
36+
public String generateHash(Object payload) {
37+
if (payload == null) return "null-payload";
38+
try {
39+
String json = (payload instanceof Object[] args)
40+
? serializeArgs(args)
41+
: hashMapper.writeValueAsString(payload);
42+
43+
MessageDigest digest = MessageDigest.getInstance("SHA-256");
44+
byte[] encodedHash = digest.digest(json.getBytes(StandardCharsets.UTF_8));
45+
return HexFormat.of().formatHex(encodedHash);
46+
} catch (Exception e) {
47+
log.warn("Hash generation failed, falling back to identity hash", e);
48+
return "fallback-" + Objects.hash(payload);
49+
}
50+
}
51+
52+
private String serializeArgs(Object[] args) throws Exception {
53+
StringBuilder sb = new StringBuilder();
54+
for (Object arg : args) {
55+
sb.append(arg == null ? "null" : hashMapper.writeValueAsString(arg));
56+
}
57+
return sb.toString();
58+
}
59+
60+
public IdempotencyRecord attemptReservation(String ns, String key, String hash, long ttl) {
61+
IdempotencyRecord initial = IdempotencyRecord.builder()
62+
.status(IdempotencyStatus.PROCESSING)
63+
.requestHash(hash)
64+
.timestamp(System.currentTimeMillis())
65+
.build();
66+
67+
Object result = store.executeLua(ns, key, initial, ttl);
68+
return (IdempotencyRecord) result;
69+
}
70+
71+
public void commit(String ns, String key, String hash, Object response, long ttl) {
72+
IdempotencyRecord completed = IdempotencyRecord.builder()
73+
.status(IdempotencyStatus.COMPLETED)
74+
.response(response)
75+
.requestHash(hash)
76+
.timestamp(System.currentTimeMillis())
77+
.build();
78+
store.save(ns, key, completed, ttl);
79+
}
80+
81+
public void rollback(String ns, String key) {
82+
store.delete(ns, key);
83+
}
84+
85+
public long getProcessingTimeoutMs() {
86+
return processingTimeoutMs;
87+
}
88+
}

0 commit comments

Comments
 (0)