Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
57 commits
Select commit Hold shift + click to select a range
61d2d07
feat: added post migrate init for log streams
nimish-ks Aug 5, 2026
9289a93
feat: logstream models
nimish-ks Aug 5, 2026
d280ec9
feat: added datadog as thid party credential provider config
nimish-ks Aug 5, 2026
46b32c3
feat: added log streams rqworker job; handle disruptions, recover gra…
nimish-ks Aug 5, 2026
7afc2d0
feat: added log stream role and update RBAC policy
nimish-ks Aug 5, 2026
eb993a9
feat: model changes and migrations
nimish-ks Aug 5, 2026
b605f66
feat: added global access user check
nimish-ks Aug 5, 2026
da59c3d
feat: added global access gating and disabled service account token a…
nimish-ks Aug 5, 2026
08a236d
feat: log stream is a enterprise tier feature
nimish-ks Aug 5, 2026
017ac79
feat: add logstream schema
nimish-ks Aug 5, 2026
38f45c0
feat: added dedicated redis queue for logstreams
nimish-ks Aug 5, 2026
ff9c828
feat: add a new url path for audit logs and wire up PublicAuditLogsView
nimish-ks Aug 5, 2026
cc37e2f
feat: include the logsteams in the third party integration disruption…
nimish-ks Aug 5, 2026
3bcdf3b
feat: chunky boy
nimish-ks Aug 5, 2026
b7781fa
feat: added the core logstreaming engine
nimish-ks Aug 5, 2026
e07f8bd
feat: handle log stream related exceptions
nimish-ks Aug 5, 2026
6afaaf7
feat: register the log sweeping job
nimish-ks Aug 5, 2026
56dce32
feat: schema for streamed logs, established - event name types, actor…
nimish-ks Aug 5, 2026
daf8073
feat: logstream sources
nimish-ks Aug 5, 2026
0e50ea2
feat: log stream adapter for destination providers
nimish-ks Aug 5, 2026
63ecdde
feat: datadod provider adapter for logstreams
nimish-ks Aug 5, 2026
7f28b7a
feat: graphql mutations for log streams
nimish-ks Aug 5, 2026
1703fca
feat: graphql query and types
nimish-ks Aug 5, 2026
f83e674
feat: quotas tests
nimish-ks Aug 5, 2026
d6dbd21
test: api paths and alterative routes
nimish-ks Aug 5, 2026
7ac9ecb
test: audit log view
nimish-ks Aug 5, 2026
e126ce1
tests: logstream engine
nimish-ks Aug 5, 2026
f965385
feat: logstream schama and types
nimish-ks Aug 5, 2026
8aa044e
feat: add logsreams in integrations
nimish-ks Aug 5, 2026
eb3f502
feat: enforce logstream RBAC
nimish-ks Aug 5, 2026
5032961
feat: lookup secret ids
nimish-ks Aug 5, 2026
9f4a0d2
feat: logstream table in audit logs and enforce global access roles
nimish-ks Aug 5, 2026
d7f5630
feat: datadog credential picker
nimish-ks Aug 5, 2026
539e786
feat: added count for disrupted integrations
nimish-ks Aug 5, 2026
e03a5d0
fix: say used with integrations instead of sync
nimish-ks Aug 5, 2026
98ccfa7
feat: added datadog icon
nimish-ks Aug 5, 2026
9ca3d7b
feat: added datadog third-party credential
nimish-ks Aug 5, 2026
6345bdc
feat: add datadog site (instance region) picker
nimish-ks Aug 5, 2026
c968901
feat: create & delete logstreams
nimish-ks Aug 5, 2026
4900c88
feat: logstream card ui
nimish-ks Aug 5, 2026
ad863b8
feat: logstream events history
nimish-ks Aug 5, 2026
22f746b
feat: logstream create
nimish-ks Aug 5, 2026
96acade
feat: logstream status
nimish-ks Aug 5, 2026
e280a44
feat: logstream
nimish-ks Aug 5, 2026
4a6359b
feat: logstream source icons
nimish-ks Aug 5, 2026
c5fb47b
feat: logstream utils
nimish-ks Aug 5, 2026
e70de7b
feat" added support for looking up secrets via their ID
nimish-ks Aug 5, 2026
d3eafee
feat: graphql queries & mutations
nimish-ks Aug 5, 2026
153e1fa
Merge branch 'main' into feat--datadog-log-stream
rohan-chaturvedi Aug 10, 2026
47ea896
fix(log-streams): remediate review findings β€” SSO alias, Retry-After …
rohan-chaturvedi Aug 12, 2026
68ded20
perf(log-streams): add concurrent partial index for unresolved delive…
rohan-chaturvedi Aug 12, 2026
ec3fec2
fix(syncing): seed Datadog site credential default and add UK1/US2-FE…
rohan-chaturvedi Aug 12, 2026
2afa258
chore(log-streams): disable audit logs REST route pending a query-per…
rohan-chaturvedi Aug 12, 2026
ec77de0
style(log-streams): align card list and create/manage dialogs with th…
rohan-chaturvedi Aug 12, 2026
a80c67b
fix(log-streams): remediate second-round review β€” backend site allowl…
rohan-chaturvedi Aug 12, 2026
5ba598f
Merge branch 'main' into feat--datadog-log-stream
rohan-chaturvedi Aug 12, 2026
6089c97
fix(log-streams): final-review remediation β€” load-more offset trackin…
rohan-chaturvedi Aug 13, 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
3 changes: 3 additions & 0 deletions .env.dev.example
Original file line number Diff line number Diff line change
Expand Up @@ -75,3 +75,6 @@ DATABASE_PASSWORD=postgres-password
REDIS_HOST=redis
REDIS_PORT=6379
REDIS_PASSWORD=

# Number of worker processes for Log Stream deliveries (optional, default 2).
#LOG_STREAM_WORKERS=2
5 changes: 5 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -68,3 +68,8 @@ DATABASE_PASSWORD=a765b221799be364c53c8a32acccf5dd90d5fc832607bdd14fccaaaa0062ad
REDIS_HOST=redis
REDIS_PORT=6379
REDIS_PASSWORD=

# Number of worker processes for Log Stream deliveries (optional, default 2).
# Log shipping is network-bound and serialized per stream, so useful
# concurrency roughly equals your number of active Log Streams.
#LOG_STREAM_WORKERS=2
9 changes: 9 additions & 0 deletions backend/api/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ def ready(self):

# Connect the post_migrate signal to a custom handler
post_migrate.connect(self.validate_licenses_post_migrate, sender=self)
post_migrate.connect(self.init_log_streams_post_migrate, sender=self)

def validate_licenses_post_migrate(self, **kwargs):

Expand All @@ -28,3 +29,11 @@ def validate_licenses_post_migrate(self, **kwargs):
activate_license(settings.PHASE_LICENSE)
except Exception as e:
logging.exception("Failed to activate license: %s", e)

def init_log_streams_post_migrate(self, **kwargs):
try:
from ee.integrations.logs.streams.jobs import init_log_stream_sweeper

init_log_stream_sweeper()
except Exception:
logging.exception("Failed to initialise log stream sweeper")
128 changes: 111 additions & 17 deletions backend/api/management/commands/rqworker.py
Original file line number Diff line number Diff line change
@@ -1,9 +1,14 @@
from django.core.management.base import BaseCommand
from django.core import management
import logging
import multiprocessing
from django.core.management.base import BaseCommand
from multiprocessing.connection import wait
import os
import sys
from django_rq.management.commands.rqworker import Command as OriginalRQWorkerCommand

logger = logging.getLogger(__name__)


class Command(BaseCommand):
help = "Runs both RQ worker and RQ scheduler in parallel"
Expand Down Expand Up @@ -50,26 +55,115 @@ def run_scheduler(self):
# to a full minute past its due time before being moved to the queue.
management.call_command("rqscheduler", "scheduled-jobs", interval=2)

def run_log_streams_workers(self):
"""Starts the worker pool for the log-streams queue.

Log shipping is network-I/O bound and per-stream serialized, so
useful concurrency ~= number of active streams. Sized via the
LOG_STREAM_WORKERS env var.

Parsed defensively: this runs in a supervised child, and a crash here
(e.g. a typo'd env value) would tear down the whole worker container
β€” syncs, emails and rotations included, not just log streams. A
zero/negative value would silently starve the queue while the
container looks healthy, so it is clamped to 1.
"""
default_workers = 2
raw = os.getenv("LOG_STREAM_WORKERS", "")
try:
workers = int(raw) if raw.strip() else default_workers
except ValueError:
logger.warning(
"Invalid LOG_STREAM_WORKERS value %r β€” using the default (%s)",
raw,
default_workers,
)
workers = default_workers
if workers < 1:
logger.warning(
"LOG_STREAM_WORKERS=%s would start no delivery workers β€” clamping to 1",
workers,
)
workers = 1
self.stdout.write(
self.style.SUCCESS(
f"Starting log-streams RQ worker pool with {workers} workers..."
)
)

management.call_command("rqworker-pool", "log-streams", num_workers=workers)

def bootstrap_log_stream_schedule(self):
"""(Re-)register the recurring log stream sweep at worker startup.

The schedule lives only in Redis. If it's lost β€” a Redis restart, or
rq-scheduler dropping an interval job whose hash expired while the
host was frozen β€” backend post_migrate wouldn't re-create it until
the next deploy. Worker startup is the natural recovery point; the
registration is idempotent (stable id, cancel-before-schedule).
"""
try:
from ee.integrations.logs.streams.jobs import init_log_stream_sweeper

init_log_stream_sweeper()
except Exception:
logger.exception("Failed to register log stream sweeper at worker startup")

def handle(self, *args, **options):
queue = options["queue"]
num_workers = options["num_workers"]

default_workers_process = multiprocessing.Process(
target=self.run_default_workers,
args=(
queue,
num_workers,
self.bootstrap_log_stream_schedule()

processes = [
multiprocessing.Process(
name="rqworker-pool-default",
target=self.run_default_workers,
args=(
queue,
num_workers,
),
),
)
scheduled_jobs_worker_process = multiprocessing.Process(
target=self.run_scheduled_jobs_worker
)
scheduler_process = multiprocessing.Process(target=self.run_scheduler)
multiprocessing.Process(
name="rqworker-scheduled-jobs",
target=self.run_scheduled_jobs_worker,
),
multiprocessing.Process(name="rqscheduler", target=self.run_scheduler),
multiprocessing.Process(
name="rqworker-pool-log-streams",
target=self.run_log_streams_workers,
),
]

for process in processes:
process.start()

default_workers_process.start()
scheduled_jobs_worker_process.start()
scheduler_process.start()
# Supervise the children instead of blindly join()ing them: a dead
# worker or scheduler process used to leave the container "Up" but
# silently degraded (e.g. rq-scheduler dying after a host sleep stops
# every recurring job with no visible failure). `wait()` blocks until
# any child's sentinel fires; exiting non-zero lets the container
# restart policy bring the whole pool back up cleanly.
try:
wait([process.sentinel for process in processes])
except KeyboardInterrupt:
self._shutdown(processes)
return

dead = next((p for p in processes if not p.is_alive()), None)
self.stderr.write(
self.style.ERROR(
f"{dead.name if dead else 'a worker process'} exited unexpectedly "
f"(exitcode={dead.exitcode if dead else '?'}); shutting down "
"worker pool for a clean restart"
)
)
self._shutdown(processes)
sys.exit(1)

default_workers_process.join()
scheduled_jobs_worker_process.join()
scheduler_process.join()
def _shutdown(self, processes):
for process in processes:
if process.is_alive():
process.terminate()
for process in processes:
process.join()
73 changes: 73 additions & 0 deletions backend/api/migrations/0132_log_streams.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
# Generated by Django 4.2.30 on 2026-08-02 14:06

from django.db import migrations, models
import django.db.models.deletion
import uuid


class Migration(migrations.Migration):

dependencies = [
('api', '0131_add_queued_sync_status'),
]

operations = [
migrations.AlterField(
model_name='auditevent',
name='resource_type',
field=models.CharField(choices=[('app', 'App'), ('env', 'Environment'), ('role', 'Role'), ('sa', 'ServiceAccount'), ('member', 'OrganisationMember'), ('policy', 'NetworkAccessPolicy'), ('pat', 'UserToken'), ('sa_token', 'ServiceAccountToken'), ('svc_token', 'ServiceToken'), ('invite', 'Invite'), ('team', 'Team'), ('rs', 'RotatingSecret'), ('stream', 'LogStream')], max_length=10),
),
migrations.AlterField(
model_name='providercredentials',
name='provider',
field=models.CharField(choices=[('cloudflare', 'Cloudflare'), ('aws', 'AWS'), ('aws_assume_role', 'AWS Assume Role'), ('github', 'GitHub'), ('gitlab', 'GitLab'), ('hashicorp_vault', 'Hashicorp Vault'), ('hashicorp_nomad', 'Hashicorp Nomad'), ('railway', 'Railway'), ('vercel', 'Vercel'), ('render', 'Render'), ('azure', 'Azure'), ('openai', 'OpenAI'), ('litellm', 'LiteLLM'), ('datadog', 'Datadog')], max_length=50),
),
migrations.CreateModel(
name='LogStream',
fields=[
('id', models.TextField(default=uuid.uuid4, editable=False, primary_key=True, serialize=False)),
('name', models.CharField(max_length=64)),
('provider', models.CharField(help_text='Log stream adapter id (resolved against the log stream adapter registry).', max_length=50)),
('sources', models.JSONField(default=list, help_text='Event source ids to ship (resolved against the log stream source registry).')),
('options', models.JSONField(default=dict)),
('max_attempts', models.PositiveIntegerField(default=5, help_text='Delivery attempts per chunk before it is recorded as failed and skipped.')),
('is_active', models.BooleanField(default=True)),
('health', models.CharField(choices=[('healthy', 'Healthy'), ('degraded', 'Degraded')], default='healthy', max_length=20)),
('paused_reason', models.TextField(blank=True, default='')),
('cursors', models.JSONField(default=dict)),
('ship_job_id', models.TextField(blank=True, null=True)),
('last_shipped_at', models.DateTimeField(blank=True, null=True)),
('last_failure_at', models.DateTimeField(blank=True, null=True)),
('last_failure_reason', models.TextField(blank=True, default='')),
('created_at', models.DateTimeField(auto_now_add=True, null=True)),
('updated_at', models.DateTimeField(auto_now=True)),
('deleted_at', models.DateTimeField(blank=True, null=True)),
('authentication', models.ForeignKey(null=True, on_delete=django.db.models.deletion.SET_NULL, related_name='log_streams', to='api.providercredentials')),
('organisation', models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name='log_streams', to='api.organisation')),
],
),
migrations.CreateModel(
name='LogStreamDeliveryEvent',
fields=[
('id', models.TextField(default=uuid.uuid4, editable=False, primary_key=True, serialize=False)),
('source', models.CharField(max_length=32)),
('status', models.CharField(choices=[('completed', 'Completed'), ('failed', 'Failed'), ('skipped', 'Skipped')], max_length=16)),
('event_count', models.PositiveIntegerField(default=0)),
('payload_bytes', models.PositiveIntegerField(default=0)),
('attempts', models.PositiveIntegerField(default=0)),
('cursor_from', models.DateTimeField(blank=True, null=True)),
('cursor_to', models.DateTimeField(blank=True, null=True)),
('resolved_at', models.DateTimeField(blank=True, null=True)),
('job_id', models.TextField(blank=True, null=True)),
('meta', models.JSONField(null=True)),
('created_at', models.DateTimeField(auto_now_add=True, null=True)),
('completed_at', models.DateTimeField(blank=True, null=True)),
('retried_from', models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name='retries', to='api.logstreamdeliveryevent')),
('stream', models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name='delivery_events', to='api.logstream')),
],
options={
'ordering': ['-created_at'],
'indexes': [models.Index(fields=['stream', '-created_at'], name='log_stream_delivery_idx')],
},
),
]
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
# Generated by Django 4.2.30 on 2026-08-03 19:35

from django.db import migrations, models


class Migration(migrations.Migration):

dependencies = [
('api', '0132_log_streams'),
]

operations = [
migrations.AddField(
model_name='logstreamdeliveryevent',
name='cursor_from_id',
field=models.TextField(blank=True, default=''),
),
migrations.AddField(
model_name='logstreamdeliveryevent',
name='cursor_to_id',
field=models.TextField(blank=True, default=''),
),
]
17 changes: 17 additions & 0 deletions backend/api/migrations/0134_logstream_drop_job_id.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
# Generated by Django 5.2.17 on 2026-08-10 10:54

from django.db import migrations


class Migration(migrations.Migration):

dependencies = [
('api', '0133_logstreamdeliveryevent_cursor_from_id_and_more'),
]

operations = [
migrations.RemoveField(
model_name='logstreamdeliveryevent',
name='job_id',
),
]
25 changes: 25 additions & 0 deletions backend/api/migrations/0135_sync_provider_service_choices.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
# Choices-only sync for EnvironmentSync.service / ProviderCredentials.provider
# (label recasing + openai/litellm/datadog additions from earlier commits).
# CharField choices are validation-level only β€” no DB DDL is emitted.

from django.db import migrations, models


class Migration(migrations.Migration):

dependencies = [
('api', '0134_logstream_drop_job_id'),
]

operations = [
migrations.AlterField(
model_name='environmentsync',
name='service',
field=models.CharField(choices=[('cloudflare_pages', 'Cloudflare Pages'), ('cloudflare_workers', 'Cloudflare Workers'), ('aws_secrets_manager', 'AWS Secrets Manager'), ('github_actions', 'GitHub Actions'), ('github_dependabot', 'GitHub Dependabot'), ('gitlab_ci', 'GitLab CI'), ('hashicorp_vault', 'HashiCorp Vault'), ('hashicorp_nomad', 'HashiCorp Nomad'), ('railway', 'Railway'), ('vercel', 'Vercel'), ('render', 'Render'), ('azure_key_vault', 'Azure Key Vault')], max_length=50),
),
migrations.AlterField(
model_name='providercredentials',
name='provider',
field=models.CharField(choices=[('cloudflare', 'Cloudflare'), ('aws', 'AWS'), ('aws_assume_role', 'AWS Assume Role'), ('github', 'GitHub'), ('gitlab', 'GitLab'), ('hashicorp_vault', 'HashiCorp Vault'), ('hashicorp_nomad', 'HashiCorp Nomad'), ('railway', 'Railway'), ('vercel', 'Vercel'), ('render', 'Render'), ('azure', 'Azure'), ('openai', 'OpenAI'), ('litellm', 'LiteLLM'), ('datadog', 'Datadog')], max_length=50),
),
]
27 changes: 27 additions & 0 deletions backend/api/migrations/0136_logstream_unresolved_idx.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
# Concurrent index build: LogStreamDeliveryEvent takes continuous writes
# from the ship path, so a plain CREATE INDEX would block them. Mirrors
# migration 0112 (SecretEvent).
#
# AddIndexConcurrently must be the ONLY operation in its (atomic=False)
# migration: if the build fails midway, any earlier operation has already
# committed and re-running the migration would crash on it, wedging the
# deploy.

from django.contrib.postgres.operations import AddIndexConcurrently
from django.db import migrations, models


class Migration(migrations.Migration):

atomic = False

dependencies = [
('api', '0135_sync_provider_service_choices'),
]

operations = [
AddIndexConcurrently(
model_name='logstreamdeliveryevent',
index=models.Index(condition=models.Q(('resolved_at__isnull', True), ('status__in', ['failed', 'skipped'])), fields=['stream', 'source'], name='log_stream_unresolved_idx'),
),
]
Loading
Loading