Skip to content

Commit 9bc64a2

Browse files
ypopovychclaude
andcommitted
Report endpoint_payload.dropped when the orchestrator removes a stored payload
`FilesOrchestrator` silently deletes stored batches that are never uploaded — too old (`maxFileAgeForRead`) or purged to keep the directory under `maxDirectorySize`. Add an `onDrop` callback (invoked with the file's byte size, outside the state lock) wired to `UploadObserver.uploadDropped`, so those removals are counted as `endpoint_payload.dropped`. Successful upload deletions (`delete(readableFile:)`) are not reported. Wired for spans and coverage. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
1 parent a6907e6 commit 9bc64a2

5 files changed

Lines changed: 109 additions & 31 deletions

File tree

Sources/DatadogSDKTesting/Telemetry/ARCHITECTURE.md

Lines changed: 9 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -144,10 +144,14 @@ methods are the protocol requirements; no-observer convenience methods are
144144

145145
### Gathered (3776)
146146
- **`endpoint_payload.*`**`requests`, `requests_ms`, `bytes` (request size),
147-
`requests_errors`, `events_count`, `events_serialization_ms`, tagged
147+
`requests_errors`, `events_count`, `events_serialization_ms`, `dropped`, tagged
148148
`test_cycle` (spans) / `code_coverage` (coverage). Wired in
149149
`DDTracer.endpointPayloadObservers(...)`. This family has **no feature call
150-
site** — the observers are its only home.
150+
site** — the observers are its only home. `dropped` is reported from
151+
`FilesOrchestrator` (via its `onDrop` callback → `UploadObserver.uploadDropped`)
152+
when a stored batch is removed without being uploaded: too old
153+
(`maxFileAgeForRead`) or purged to keep the directory under `maxDirectorySize`.
154+
Successful upload deletions (`delete(readableFile:)`) are **not** drops.
151155
- **API request families**`git_requests.{settings,search_commits,objects_pack}`
152156
(+ `_ms`, `_errors`, and `objects_pack_bytes`), `itr_skippable_tests.{request,
153157
request_ms,request_errors,response_bytes}`, `known_tests.{request,request_ms,
@@ -164,12 +168,7 @@ methods are the protocol requirements; no-observer convenience methods are
164168
- `itr_skippable_tests.response_tests` / `response_suites`
165169
- `known_tests.response_tests`
166170
- `test_management_tests.response_tests`
167-
2. **`endpoint_payload.dropped`** — the `UploadObserver.uploadDropped` hook exists
168-
but the worker has no retry-exhaustion drop today; failed batches stay on disk
169-
and are age-purged by `FilesOrchestrator`. Wire `dropped` from the purge path
170-
(or add an explicit drop) — see `DDTracer.endpointPayloadObservers` where
171-
`onDropped` is already mapped to `endpointPayload.dropped`.
172-
3. **Local feature metrics** — emitted directly via `SessionConfig.telemetry` /
171+
2. **Local feature metrics** — emitted directly via `SessionConfig.telemetry` /
173172
the feature's injected `Telemetry`, at the feature instrumentation sites:
174173
- `events.created` / `events.finished` (+ all their tags: `event_type`,
175174
`test_framework`, `is_new`, `is_modified`, `is_retry`, `retry_reason`,
@@ -182,9 +181,9 @@ methods are the protocol requirements; no-observer convenience methods are
182181
- `git.commit_sha_match` / `git.commit_sha_discrepancy` — git info providers.
183182
- `itr.skipped` / `itr.unskippable` / `itr.forced_run``TestImpactAnalysis`.
184183
- `code_coverage.{started,finished,is_empty,errors,files}` — coverage feature.
185-
4. **`impacted_tests_detection.*`** — no feature/API exists yet; instruments are
184+
3. **`impacted_tests_detection.*`** — no feature/API exists yet; instruments are
186185
defined but unused.
187-
5. **`error_type` granularity**`Telemetry.errorType(statusCode:)` maps `nil`
186+
4. **`error_type` granularity**`Telemetry.errorType(statusCode:)` maps `nil`
188187
status to `.network`; it cannot distinguish `timeout` from `network` without
189188
the underlying `URLError`. Refine if needed.
190189

Sources/EventsExporter/CoverageExporter/CoverageExporter.swift

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -40,10 +40,12 @@ internal final class CoverageExporter: CoverageExporterType {
4040
observers: ExporterObservers.Feature = .init()) throws {
4141
self.configuration = config
4242

43+
let uploadObserver = observers.upload
4344
let filesOrchestrator = FilesOrchestrator(
4445
directory: try storage.createSubdirectory(path: "v1"),
4546
performance: PerformancePreset.instantDataDelivery,
46-
dateProvider: SystemDateProvider()
47+
dateProvider: SystemDateProvider(),
48+
onDrop: uploadObserver.map { obs in { obs.uploadDropped(payloadBytes: $0) } }
4749
)
4850

4951
let encoder = api.encoder

Sources/EventsExporter/Persistence/FilesOrchestrator.swift

Lines changed: 53 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -148,47 +148,55 @@ internal final class FilesOrchestrator: FilesOrchestratorType {
148148

149149
// MARK: - Reader
150150

151+
/// Reader results paired with the byte sizes of any files dropped
152+
/// (deleted for exceeding `maxFileAgeForRead`) during the scan. The
153+
/// caller reports the drops *outside* the state lock.
151154
func oldestReadableFile(directory: borrowing Directory,
152155
performance: borrowing StoragePerformancePreset,
153-
dateProvider: borrowing DateProvider) throws -> ReadableFile?
156+
dateProvider: borrowing DateProvider) throws -> (file: ReadableFile?, droppedBytes: [Int])
154157
{
155-
guard let oldest = try fileInfos(directory: directory,
156-
performance: performance,
157-
dateProvider: dateProvider).first
158-
else { return nil }
158+
let (infos, droppedBytes) = try fileInfos(directory: directory,
159+
performance: performance,
160+
dateProvider: dateProvider)
161+
guard let oldest = infos.first else { return (nil, droppedBytes) }
159162
let age = dateProvider.currentDate().timeIntervalSince(oldest.creationDate)
160-
return age >= performance.minFileAgeForRead ? oldest.file : nil
163+
return (age >= performance.minFileAgeForRead ? oldest.file : nil, droppedBytes)
161164
}
162165

163166
func allReadableFiles(directory: borrowing Directory,
164167
performance: borrowing StoragePerformancePreset,
165-
dateProvider: borrowing DateProvider) throws -> [ReadableFile]
168+
dateProvider: borrowing DateProvider) throws -> (files: [ReadableFile], droppedBytes: [Int])
166169
{
167-
try fileInfos(directory: directory,
168-
performance: performance,
169-
dateProvider: dateProvider).map { $0.file }
170+
let (infos, droppedBytes) = try fileInfos(directory: directory,
171+
performance: performance,
172+
dateProvider: dateProvider)
173+
return (infos.map { $0.file }, droppedBytes)
170174
}
171175

172176
private func fileInfos(directory: borrowing Directory,
173177
performance: borrowing StoragePerformancePreset,
174-
dateProvider: borrowing DateProvider) throws -> [FileInfo]
178+
dateProvider: borrowing DateProvider) throws -> (files: [FileInfo], droppedBytes: [Int])
175179
{
176180
let allFiles = try directory.files()
177181
.filter { !activeWrites.contains($0.name) }
178182
.map { FileInfo(file: $0) }
179183

180184
var readableFiles: [FileInfo] = []
181185
readableFiles.reserveCapacity(allFiles.count)
186+
var droppedBytes: [Int] = []
182187
for info in allFiles {
183188
let fileAge = dateProvider.currentDate().timeIntervalSince(info.creationDate)
184189
if fileAge > performance.maxFileAgeForRead {
190+
// Too old to ever upload — count it as a dropped payload.
191+
let size = (try? info.file.size()).map(Int.init) ?? 0
185192
try info.file.delete()
193+
droppedBytes.append(size)
186194
} else {
187195
readableFiles.append(info)
188196
}
189197
}
190198

191-
return readableFiles.sorted()
199+
return (readableFiles.sorted(), droppedBytes)
192200
}
193201
}
194202

@@ -200,14 +208,20 @@ internal final class FilesOrchestrator: FilesOrchestratorType {
200208
private let directory: Directory
201209
private let dateProvider: DateProvider
202210
private let performance: StoragePerformancePreset
211+
/// Invoked (outside the state lock) with the byte size of each file removed
212+
/// without being uploaded — too old (`maxFileAgeForRead`) or purged to keep
213+
/// the directory under `maxDirectorySize`. Wired to `endpoint_payload.dropped`.
214+
private let onDrop: (@Sendable (Int) -> Void)?
203215

204216
init(directory: Directory,
205217
performance: StoragePerformancePreset,
206-
dateProvider: DateProvider)
218+
dateProvider: DateProvider,
219+
onDrop: (@Sendable (Int) -> Void)? = nil)
207220
{
208221
self.directory = directory
209222
self.dateProvider = dateProvider
210223
self.performance = performance
224+
self.onDrop = onDrop
211225
self.state = Synced(.init())
212226
}
213227

@@ -245,15 +259,30 @@ internal final class FilesOrchestrator: FilesOrchestratorType {
245259
// MARK: - `ReadableFile` orchestration
246260

247261
func getReadableFile() throws -> ReadableFile? {
248-
try state.use { try $0.oldestReadableFile(directory: directory,
249-
performance: performance,
250-
dateProvider: dateProvider) }
262+
let (file, droppedBytes) = try state.use {
263+
try $0.oldestReadableFile(directory: directory,
264+
performance: performance,
265+
dateProvider: dateProvider)
266+
}
267+
reportDrops(droppedBytes)
268+
return file
251269
}
252270

253271
func getAllReadableFiles() throws -> [ReadableFile] {
254-
try state.use { try $0.allReadableFiles(directory: directory,
255-
performance: performance,
256-
dateProvider: dateProvider) }
272+
let (files, droppedBytes) = try state.use {
273+
try $0.allReadableFiles(directory: directory,
274+
performance: performance,
275+
dateProvider: dateProvider)
276+
}
277+
reportDrops(droppedBytes)
278+
return files
279+
}
280+
281+
/// Report dropped-payload sizes. Called outside the state lock so the
282+
/// observer (which may touch its own locks) can't contend with file ops.
283+
private func reportDrops(_ droppedBytes: [Int]) {
284+
guard let onDrop else { return }
285+
droppedBytes.forEach(onDrop)
257286
}
258287

259288
func delete(readableFile: ReadableFile) throws {
@@ -272,7 +301,10 @@ internal final class FilesOrchestrator: FilesOrchestratorType {
272301
.compactMap { (info) -> FileInfo? in
273302
let fileAge = dateProvider.currentDate().timeIntervalSince(info.creationDate)
274303
if fileAge > performance.maxFileAgeForRead {
304+
// Too old to ever upload — count it as a dropped payload.
305+
let size = (try? info.file.size()).map(Int.init) ?? 0
275306
try info.file.delete()
307+
onDrop?(size)
276308
return nil
277309
}
278310
return info
@@ -286,6 +318,8 @@ internal final class FilesOrchestrator: FilesOrchestratorType {
286318
let fileWithSize = filesWithSize.removeFirst()
287319
try fileWithSize.file.delete()
288320
sizeFreed += fileWithSize.size
321+
// Purged to stay under the directory size limit — also a drop.
322+
onDrop?(Int(fileWithSize.size))
289323
}
290324
}
291325
}

Sources/EventsExporter/Spans/SpansExporter.swift

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,10 +17,12 @@ internal final class SpansExporter: SpanExporter {
1717
observers: ExporterObservers.Feature = .init()) throws {
1818
self.configuration = config
1919

20+
let uploadObserver = observers.upload
2021
let filesOrchestrator = FilesOrchestrator(
2122
directory: try storage.createSubdirectory(path: "v1"),
2223
performance: configuration.performancePreset,
23-
dateProvider: SystemDateProvider()
24+
dateProvider: SystemDateProvider(),
25+
onDrop: uploadObserver.map { obs in { obs.uploadDropped(payloadBytes: $0) } }
2426
)
2527

2628
var metadata = SpanSanitizer().sanitize(metadata: config.metadata)

Tests/EventsExporter/Persistence/FilesOrchestratorTests.swift

Lines changed: 41 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -223,6 +223,47 @@ class FilesOrchestratorTests: XCTestCase {
223223
XCTAssertEqual(try temporaryDirectory.files().count, 0)
224224
}
225225

226+
func testGivenFileTooOld_whenScanned_itReportsDropWithSize() throws {
227+
final class Box: @unchecked Sendable { var sizes: [Int] = [] }
228+
let box = Box()
229+
let dateProvider = RelativeDateProvider()
230+
let orchestrator = FilesOrchestrator(
231+
directory: temporaryDirectory,
232+
performance: performance,
233+
dateProvider: dateProvider,
234+
onDrop: { box.sizes.append($0) }
235+
)
236+
let file = try temporaryDirectory.createFile(named: dateProvider.currentDate().toFileName)
237+
try file.append(data: Data("hello".utf8)) // 5 bytes
238+
239+
dateProvider.advance(bySeconds: 2 * performance.maxFileAgeForRead)
240+
241+
// Scanning for a readable file finds it too old, deletes it, and reports
242+
// the drop with the file's size.
243+
XCTAssertNil(try orchestrator.getReadableFile())
244+
XCTAssertEqual(try temporaryDirectory.files().count, 0)
245+
XCTAssertEqual(box.sizes, [5])
246+
}
247+
248+
func testGivenSuccessfulRead_whenDeleted_itDoesNotReportDrop() throws {
249+
final class Box: @unchecked Sendable { var sizes: [Int] = [] }
250+
let box = Box()
251+
let dateProvider = RelativeDateProvider()
252+
let orchestrator = FilesOrchestrator(
253+
directory: temporaryDirectory,
254+
performance: performance,
255+
dateProvider: dateProvider,
256+
onDrop: { box.sizes.append($0) }
257+
)
258+
_ = try temporaryDirectory.createFile(named: dateProvider.currentDate().toFileName)
259+
dateProvider.advance(bySeconds: 1 + performance.minFileAgeForRead)
260+
261+
let readableFile = try orchestrator.getReadableFile().unwrapOrThrow()
262+
try orchestrator.delete(readableFile: readableFile) // marked-as-read, not a drop
263+
264+
XCTAssertEqual(box.sizes, [], "Successful upload deletion must not be reported as a drop")
265+
}
266+
226267
func testItDeletesReadableFile() throws {
227268
let dateProvider = RelativeDateProvider()
228269
let orchestrator = configureOrchestrator(using: dateProvider)

0 commit comments

Comments
 (0)