Skip to content

Commit c98ef4d

Browse files
committed
feat(api): expose read-only yarn application metadata
1 parent 785062e commit c98ef4d

7 files changed

Lines changed: 407 additions & 2 deletions

File tree

README.md

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -148,6 +148,18 @@ application/container information before an `ERROR` or `WARNING` event is sent
148148
to the client. Log contents and SPNEGO cookies are never written by these
149149
diagnostic messages.
150150

151+
The application image and starter also expose read-only application metadata
152+
through the same `YarnLogAuthorizer` used by log streams:
153+
154+
```text
155+
GET /api/v1/yarn/applications/{applicationId}
156+
GET /api/v1/yarn/applications/{applicationId}/attempts
157+
GET /api/v1/yarn/applications/{applicationId}/containers
158+
```
159+
160+
These endpoints return API-owned JSON models rather than Hadoop implementation
161+
objects. Mutating YARN operations are intentionally not exposed.
162+
151163
The image classpath contains `/etc/hadoop/conf` and `/app/extensions/*`. Mount
152164
extra JARs into the extension directory when adding a filesystem provider,
153165
Spring auto-configuration, metrics exporter, or another runtime integration:
Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,46 @@
1+
package org.openprojectx.hadoop.yarn.log.api.autoconfigure
2+
3+
import org.openprojectx.hadoop.yarn.log.api.YarnApplicationAccessDeniedException
4+
import org.openprojectx.hadoop.yarn.log.api.YarnApplicationAttemptInfo
5+
import org.openprojectx.hadoop.yarn.log.api.YarnApplicationInfo
6+
import org.openprojectx.hadoop.yarn.log.api.YarnApplicationQueryService
7+
import org.openprojectx.hadoop.yarn.log.api.YarnContainerInfo
8+
import org.springframework.http.HttpStatus
9+
import org.springframework.web.bind.annotation.GetMapping
10+
import org.springframework.web.bind.annotation.PathVariable
11+
import org.springframework.web.bind.annotation.RestController
12+
import org.springframework.web.server.ResponseStatusException
13+
import reactor.core.publisher.Flux
14+
import reactor.core.publisher.Mono
15+
import java.security.Principal
16+
17+
@RestController
18+
class YarnApplicationController(
19+
private val service: YarnApplicationQueryService,
20+
) {
21+
@GetMapping("/api/v1/yarn/applications/{applicationId}")
22+
fun application(
23+
@PathVariable("applicationId") applicationId: String,
24+
principal: Principal?,
25+
): Mono<YarnApplicationInfo> = service.application(principal?.name, applicationId).mapAccessDenied()
26+
27+
@GetMapping("/api/v1/yarn/applications/{applicationId}/attempts")
28+
fun attempts(
29+
@PathVariable("applicationId") applicationId: String,
30+
principal: Principal?,
31+
): Flux<YarnApplicationAttemptInfo> = service.attempts(principal?.name, applicationId).mapAccessDenied()
32+
33+
@GetMapping("/api/v1/yarn/applications/{applicationId}/containers")
34+
fun containers(
35+
@PathVariable("applicationId") applicationId: String,
36+
principal: Principal?,
37+
): Flux<YarnContainerInfo> = service.containers(principal?.name, applicationId).mapAccessDenied()
38+
39+
private fun <T : Any> Mono<T>.mapAccessDenied(): Mono<T> = onErrorMap(YarnApplicationAccessDeniedException::class.java) {
40+
ResponseStatusException(HttpStatus.FORBIDDEN, it.message, it)
41+
}
42+
43+
private fun <T : Any> Flux<T>.mapAccessDenied(): Flux<T> = onErrorMap(YarnApplicationAccessDeniedException::class.java) {
44+
ResponseStatusException(HttpStatus.FORBIDDEN, it.message, it)
45+
}
46+
}

autoconfigure/src/main/kotlin/org/openprojectx/hadoop/yarn/log/api/autoconfigure/YarnLogApiAutoConfiguration.kt

Lines changed: 22 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -5,9 +5,11 @@ import org.apache.hadoop.conf.Configuration
55
import org.apache.hadoop.yarn.client.api.YarnClient
66
import org.apache.hadoop.yarn.conf.YarnConfiguration
77
import org.apache.hadoop.yarn.webapp.util.WebAppUtils
8+
import org.openprojectx.hadoop.yarn.log.api.YarnApplicationQueryService
89
import org.openprojectx.hadoop.yarn.log.api.YarnLogAuthorizer
910
import org.openprojectx.hadoop.yarn.log.api.YarnLogStreamService
1011
import org.openprojectx.hadoop.yarn.log.api.engine.AggregatedLogSource
12+
import org.openprojectx.hadoop.yarn.log.api.engine.DefaultYarnApplicationQueryService
1113
import org.openprojectx.hadoop.yarn.log.api.engine.DefaultYarnLogStreamService
1214
import org.openprojectx.hadoop.yarn.log.api.engine.HadoopSpnegoCookieProvider
1315
import org.openprojectx.hadoop.yarn.log.api.engine.HadoopYarnGateway
@@ -64,11 +66,25 @@ class YarnLogApiAutoConfiguration {
6466
@ConditionalOnMissingBean
6567
fun yarnLogAuthorizer(): YarnLogAuthorizer = YarnLogAuthorizer { _, _, _ -> reactor.core.publisher.Mono.just(true) }
6668

69+
@Bean
70+
@ConditionalOnMissingBean
71+
fun hadoopYarnGateway(
72+
yarnClient: YarnClient,
73+
@Qualifier("yarnLogHadoopScheduler") scheduler: Scheduler,
74+
) = HadoopYarnGateway(yarnClient, scheduler)
75+
76+
@Bean
77+
@ConditionalOnMissingBean
78+
fun yarnApplicationQueryService(
79+
yarnGateway: HadoopYarnGateway,
80+
authorizer: YarnLogAuthorizer,
81+
): YarnApplicationQueryService = DefaultYarnApplicationQueryService(yarnGateway, authorizer)
82+
6783
@Bean
6884
@ConditionalOnMissingBean
6985
fun yarnLogStreamService(
7086
configuration: Configuration,
71-
yarnClient: YarnClient,
87+
yarnGateway: HadoopYarnGateway,
7288
webClientBuilder: WebClient.Builder,
7389
@Qualifier("yarnLogHadoopScheduler") scheduler: Scheduler,
7490
authorizer: YarnLogAuthorizer,
@@ -89,7 +105,7 @@ class YarnLogApiAutoConfiguration {
89105
.build()
90106
val cookieProvider = HadoopSpnegoCookieProvider(scheduler, properties.spnegoCookieTtl)
91107
return DefaultYarnLogStreamService(
92-
yarnGateway = HadoopYarnGateway(yarnClient, scheduler),
108+
yarnGateway = yarnGateway,
93109
nodeManagerClient = NodeManagerLogClient(
94110
webClient,
95111
cookieProvider,
@@ -106,6 +122,10 @@ class YarnLogApiAutoConfiguration {
106122
fun yarnLogSseController(service: YarnLogStreamService, properties: YarnLogApiProperties) =
107123
YarnLogSseController(service, properties)
108124

125+
@Bean
126+
@ConditionalOnMissingBean
127+
fun yarnApplicationController(service: YarnApplicationQueryService) = YarnApplicationController(service)
128+
109129
@Bean
110130
@ConditionalOnMissingBean(YarnLogWebSocketHandler::class)
111131
fun yarnLogWebSocketHandler(
Lines changed: 127 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,127 @@
1+
package org.openprojectx.hadoop.yarn.log.api.autoconfigure
2+
3+
import com.ninjasquad.springmockk.MockkBean
4+
import io.mockk.every
5+
import org.openprojectx.hadoop.yarn.log.api.YarnApplicationAccessDeniedException
6+
import org.openprojectx.hadoop.yarn.log.api.YarnApplicationAttemptInfo
7+
import org.openprojectx.hadoop.yarn.log.api.YarnApplicationInfo
8+
import org.openprojectx.hadoop.yarn.log.api.YarnApplicationQueryService
9+
import org.openprojectx.hadoop.yarn.log.api.YarnContainerInfo
10+
import org.springframework.beans.factory.annotation.Autowired
11+
import org.springframework.boot.autoconfigure.EnableAutoConfiguration
12+
import org.springframework.boot.webflux.test.autoconfigure.WebFluxTest
13+
import org.springframework.context.annotation.Configuration
14+
import org.springframework.test.context.ContextConfiguration
15+
import org.springframework.test.web.reactive.server.WebTestClient
16+
import reactor.core.publisher.Flux
17+
import reactor.core.publisher.Mono
18+
import kotlin.test.Test
19+
20+
@WebFluxTest(YarnApplicationController::class)
21+
@ContextConfiguration(
22+
classes = [YarnApplicationControllerTest.TestApplication::class, YarnApplicationController::class],
23+
)
24+
class YarnApplicationControllerTest {
25+
@Autowired
26+
private lateinit var webTestClient: WebTestClient
27+
28+
@MockkBean
29+
private lateinit var service: YarnApplicationQueryService
30+
31+
@Test
32+
fun `returns application attempts and containers as API models`() {
33+
every { service.application(null, APPLICATION_ID) } returns Mono.just(application())
34+
every { service.attempts(null, APPLICATION_ID) } returns Flux.just(attempt())
35+
every { service.containers(null, APPLICATION_ID) } returns Flux.just(container())
36+
37+
webTestClient.get().uri("/api/v1/yarn/applications/$APPLICATION_ID")
38+
.exchange()
39+
.expectStatus().isOk
40+
.expectBody()
41+
.jsonPath("$.applicationId").isEqualTo(APPLICATION_ID)
42+
.jsonPath("$.state").isEqualTo("RUNNING")
43+
.jsonPath("$.user").isEqualTo("alice")
44+
45+
webTestClient.get().uri("/api/v1/yarn/applications/$APPLICATION_ID/attempts")
46+
.exchange()
47+
.expectStatus().isOk
48+
.expectBody()
49+
.jsonPath("$[0].attemptId").isEqualTo(ATTEMPT_ID)
50+
51+
webTestClient.get().uri("/api/v1/yarn/applications/$APPLICATION_ID/containers")
52+
.exchange()
53+
.expectStatus().isOk
54+
.expectBody()
55+
.jsonPath("$[0].containerId").isEqualTo(CONTAINER_ID)
56+
.jsonPath("$[0].memoryMb").isEqualTo(1024)
57+
}
58+
59+
@Test
60+
fun `maps authorization denial to forbidden`() {
61+
every { service.application(null, APPLICATION_ID) } returns
62+
Mono.error(YarnApplicationAccessDeniedException(APPLICATION_ID))
63+
64+
webTestClient.get().uri("/api/v1/yarn/applications/$APPLICATION_ID")
65+
.exchange()
66+
.expectStatus().isForbidden
67+
}
68+
69+
private fun application() = YarnApplicationInfo(
70+
applicationId = APPLICATION_ID,
71+
currentAttemptId = ATTEMPT_ID,
72+
name = "example",
73+
applicationType = "SPARK",
74+
user = "alice",
75+
queue = "default",
76+
state = "RUNNING",
77+
finalStatus = "UNDEFINED",
78+
progress = 0.5f,
79+
trackingUrl = "http://rm/proxy/$APPLICATION_ID",
80+
originalTrackingUrl = null,
81+
diagnostics = null,
82+
submitTime = 1,
83+
startTime = 2,
84+
launchTime = 3,
85+
finishTime = 0,
86+
logAggregationStatus = "RUNNING_WITH_FAILURE",
87+
tags = setOf("team=data"),
88+
)
89+
90+
private fun attempt() = YarnApplicationAttemptInfo(
91+
attemptId = ATTEMPT_ID,
92+
state = "RUNNING",
93+
amContainerId = CONTAINER_ID,
94+
host = "worker.example.com",
95+
rpcPort = 8042,
96+
trackingUrl = null,
97+
originalTrackingUrl = null,
98+
diagnostics = null,
99+
startTime = 2,
100+
finishTime = 0,
101+
)
102+
103+
private fun container() = YarnContainerInfo(
104+
containerId = CONTAINER_ID,
105+
attemptId = ATTEMPT_ID,
106+
state = "RUNNING",
107+
nodeId = "worker.example.com:45454",
108+
nodeHttpAddress = "worker.example.com:8042",
109+
logUrl = null,
110+
diagnostics = null,
111+
exitStatus = -1000,
112+
memoryMb = 1024,
113+
virtualCores = 1,
114+
creationTime = 3,
115+
finishTime = 0,
116+
)
117+
118+
private companion object {
119+
const val APPLICATION_ID = "application_123_0001"
120+
const val ATTEMPT_ID = "appattempt_123_0001_000001"
121+
const val CONTAINER_ID = "container_123_0001_01_000001"
122+
}
123+
124+
@Configuration(proxyBeanMethods = false)
125+
@EnableAutoConfiguration
126+
class TestApplication
127+
}
Lines changed: 64 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,64 @@
1+
package org.openprojectx.hadoop.yarn.log.api
2+
3+
import reactor.core.publisher.Flux
4+
import reactor.core.publisher.Mono
5+
6+
interface YarnApplicationQueryService {
7+
fun application(requester: String?, applicationId: String): Mono<YarnApplicationInfo>
8+
9+
fun attempts(requester: String?, applicationId: String): Flux<YarnApplicationAttemptInfo>
10+
11+
fun containers(requester: String?, applicationId: String): Flux<YarnContainerInfo>
12+
}
13+
14+
data class YarnApplicationInfo(
15+
val applicationId: String,
16+
val currentAttemptId: String?,
17+
val name: String,
18+
val applicationType: String,
19+
val user: String,
20+
val queue: String,
21+
val state: String,
22+
val finalStatus: String,
23+
val progress: Float,
24+
val trackingUrl: String?,
25+
val originalTrackingUrl: String?,
26+
val diagnostics: String?,
27+
val submitTime: Long,
28+
val startTime: Long,
29+
val launchTime: Long,
30+
val finishTime: Long,
31+
val logAggregationStatus: String?,
32+
val tags: Set<String>,
33+
)
34+
35+
data class YarnApplicationAttemptInfo(
36+
val attemptId: String,
37+
val state: String,
38+
val amContainerId: String?,
39+
val host: String?,
40+
val rpcPort: Int,
41+
val trackingUrl: String?,
42+
val originalTrackingUrl: String?,
43+
val diagnostics: String?,
44+
val startTime: Long,
45+
val finishTime: Long,
46+
)
47+
48+
data class YarnContainerInfo(
49+
val containerId: String,
50+
val attemptId: String,
51+
val state: String,
52+
val nodeId: String?,
53+
val nodeHttpAddress: String?,
54+
val logUrl: String?,
55+
val diagnostics: String?,
56+
val exitStatus: Int,
57+
val memoryMb: Long?,
58+
val virtualCores: Int?,
59+
val creationTime: Long,
60+
val finishTime: Long,
61+
)
62+
63+
class YarnApplicationAccessDeniedException(applicationId: String) :
64+
SecurityException("Not authorized to read $applicationId")
Lines changed: 65 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,65 @@
1+
package org.openprojectx.hadoop.yarn.log.api.engine
2+
3+
import org.openprojectx.hadoop.yarn.log.api.YarnApplicationAccessDeniedException
4+
import org.openprojectx.hadoop.yarn.log.api.YarnApplicationAttemptInfo
5+
import org.openprojectx.hadoop.yarn.log.api.YarnApplicationInfo
6+
import org.openprojectx.hadoop.yarn.log.api.YarnApplicationQueryService
7+
import org.openprojectx.hadoop.yarn.log.api.YarnContainerInfo
8+
import org.openprojectx.hadoop.yarn.log.api.YarnLogAuthorizer
9+
import org.slf4j.LoggerFactory
10+
import reactor.core.publisher.Flux
11+
import reactor.core.publisher.Mono
12+
13+
class DefaultYarnApplicationQueryService(
14+
private val yarnGateway: HadoopYarnGateway,
15+
private val authorizer: YarnLogAuthorizer,
16+
) : YarnApplicationQueryService {
17+
override fun application(requester: String?, applicationId: String): Mono<YarnApplicationInfo> =
18+
authorizedApplication(requester, applicationId)
19+
.doOnSubscribe { logQuery("application", applicationId, requester) }
20+
.doOnError { error -> logFailure("application", applicationId, requester, error) }
21+
22+
override fun attempts(requester: String?, applicationId: String): Flux<YarnApplicationAttemptInfo> =
23+
authorizedApplication(requester, applicationId)
24+
.flatMapMany { yarnGateway.attempts(applicationId) }
25+
.doOnSubscribe { logQuery("attempts", applicationId, requester) }
26+
.doOnError { error -> logFailure("attempts", applicationId, requester, error) }
27+
28+
override fun containers(requester: String?, applicationId: String): Flux<YarnContainerInfo> =
29+
authorizedApplication(requester, applicationId)
30+
.flatMapMany { yarnGateway.containers(applicationId) }
31+
.doOnSubscribe { logQuery("containers", applicationId, requester) }
32+
.doOnError { error -> logFailure("containers", applicationId, requester, error) }
33+
34+
private fun authorizedApplication(requester: String?, applicationId: String): Mono<YarnApplicationInfo> =
35+
yarnGateway.application(applicationId).flatMap { application ->
36+
authorizer.authorize(requester, application.applicationId, application.user)
37+
.flatMap { allowed ->
38+
if (allowed) Mono.just(application)
39+
else Mono.error(YarnApplicationAccessDeniedException(applicationId))
40+
}
41+
}
42+
43+
private fun logQuery(operation: String, applicationId: String, requester: String?) {
44+
logger.info(
45+
"Querying YARN application: operation={}, applicationId={}, requester={}",
46+
operation,
47+
applicationId,
48+
requester ?: "anonymous",
49+
)
50+
}
51+
52+
private fun logFailure(operation: String, applicationId: String, requester: String?, error: Throwable) {
53+
logger.error(
54+
"YARN application query failed: operation={}, applicationId={}, requester={}",
55+
operation,
56+
applicationId,
57+
requester ?: "anonymous",
58+
error,
59+
)
60+
}
61+
62+
private companion object {
63+
val logger = LoggerFactory.getLogger(DefaultYarnApplicationQueryService::class.java)
64+
}
65+
}

0 commit comments

Comments
 (0)