De 87 referansefilene bar en plain-text `| Verified: <dato>`-hale på **Last updated:**-linjen i 500B-header-vinduet — usynlig for den bold-only kontrakt-stacken (kb-headers.mjs / audit RE_VERIFIED), og claimet en verifisering judgen aldri gjorde (samme poison-klasse som de 14 bold **Verified:** MCP Spor 1 fjernet). Uhåndtert springer den også dual-Verified-fellen: R7s insertVerifiedFields ville stemplet en bold-verdi ved siden av den plain → to motstridende provenance-claims per fil. - ny driver strip-stale-verified-pipe.mjs: frosset 87-manifest (18 advisor + 45 eng + 8 gov + 16 sec), pure verdi-bevarende strip (kun ` | Verified: …`-halen; **Last updated:**-dato byte-eksakt), hard per-fil-invariant (linjeantall uendret, body byte-identisk, dato bevart), idempotent, atomicWriteSync (RX-OPS2 recovery-kontrakt). - audit-corpus-headers.mjs: ny plain-Verified-deteksjon (RE_PLAIN_VERIFIED + plainVerifiedPipe) — gjør M4-blindheten synlig så en stale plain-hale ikke kan gjenoppstå stille (non-advisor scope). - 87 filer strippet; plain Verified i vinduet 0/389; live-audit plainVerifiedPipe 0. Mekanisme: +15 tester (12 strip + 3 audit). Suite 875→890 exit 0. validate-plugin.sh 250/0. Utsatt → RX-KB1b: footer-dato-avvik + label-whitelist (annen dialekt, flag-to-human).
20 KiB
Asynchronous Processing Patterns
Last updated: 2026-06-24 Status: GA Category: Performance & Scalability Type: reference Source: https://learn.microsoft.com/azure/architecture/guide/architecture-styles/event-driven
Innhold
- Introduksjon
- Kjernekomponenter
- Queue-based Architectures
- Event-Driven Design
- Request-Response Decoupling
- Status Polling and Webhooks
- Event-Driven Architecture Styles (oppdatert 2026-04)
- Norsk offentlig sektor
- Beslutningsrammeverk
- Referanser
- For Cosmo
Introduksjon
Asynkron prosessering er en arkitekturstrategi der AI-forespørsler behandles uavhengig av den opprinnelige klientforbindelsen. I stedet for at klienten venter synkront på et svar fra Azure OpenAI (som kan ta fra 500ms til flere minutter for reasoning-modeller), plasseres forespørselen i en kø, behandles i bakgrunnen, og resultatet leveres via polling, webhook eller push-notifikasjon.
For Azure OpenAI tilbyr Microsoft flere innebygde asynkrone mekanismer: Batch API for store volum, Background Tasks i Responses API for langvarige oppgaver, og Webhooks for hendelsesbasert leveranse. I tillegg kan organisasjoner bygge egne asynkrone arkitekturer med Azure Service Bus, Azure Queue Storage eller Azure Event Hubs som mellomlag.
I norsk offentlig sektor er asynkron prosessering spesielt relevant for dokumentanalyse, saksbehandlingsstøtte og rapportgenerering — oppgaver der brukeren ikke trenger umiddelbart svar, men der volumet kan være svært høyt i perioder (f.eks. ved frister for høringssvar eller klagebehandling).
Kjernekomponenter
| Komponent | Formål | Teknologi |
|---|---|---|
| Azure Service Bus | Enterprise message broker med køer og topics | Azure Service Bus |
| Azure Queue Storage | Enkel, kostnadseffektiv meldingskø | Azure Storage |
| Azure Event Hubs | Høy-throughput event streaming | Azure Event Hubs |
| Azure Functions | Serverless compute for kø-triggered prosessering | Azure Functions |
| Batch API | Innebygd asynkron batch-prosessering | Azure OpenAI |
| Background Tasks | Langvarige oppgaver i Responses API | Azure OpenAI |
| Webhooks | Hendelsesbasert notifikasjon | Azure OpenAI |
Queue-based Architectures
Service Bus-basert AI-prosessering
# Producer: Legg forespørsler i kø
from azure.servicebus import ServiceBusClient, ServiceBusMessage
import json
class AIRequestProducer:
"""Queue AI requests via Azure Service Bus."""
def __init__(self, connection_string: str, queue_name: str = "ai-requests"):
self.client = ServiceBusClient.from_connection_string(connection_string)
self.sender = self.client.get_queue_sender(queue_name)
async def submit_request(
self,
request_id: str,
messages: list[dict],
priority: str = "normal",
callback_url: str = None
) -> str:
"""Submit AI request to queue. Returns request ID for polling."""
payload = {
"request_id": request_id,
"messages": messages,
"priority": priority,
"callback_url": callback_url,
"submitted_at": datetime.utcnow().isoformat()
}
message = ServiceBusMessage(
body=json.dumps(payload),
message_id=request_id,
subject=priority,
session_id=request_id if priority == "urgent" else None,
time_to_live=timedelta(hours=24)
)
await self.sender.send_messages(message)
return request_id
# Consumer: Prosesser forespørsler fra kø
from azure.servicebus.aio import ServiceBusClient as AsyncServiceBusClient
from openai import AsyncAzureOpenAI
class AIRequestConsumer:
"""Process AI requests from Service Bus queue."""
def __init__(
self,
sb_connection: str,
queue_name: str,
openai_client: AsyncAzureOpenAI,
max_concurrent: int = 10
):
self.sb_client = AsyncServiceBusClient.from_connection_string(
sb_connection)
self.queue_name = queue_name
self.openai = openai_client
self.semaphore = asyncio.Semaphore(max_concurrent)
async def process_messages(self):
"""Continuously process messages from queue."""
async with self.sb_client.get_queue_receiver(
self.queue_name,
max_wait_time=30
) as receiver:
async for message in receiver:
asyncio.create_task(
self._handle_message(receiver, message))
async def _handle_message(self, receiver, message):
async with self.semaphore:
try:
payload = json.loads(str(message))
# Prosesser med Azure OpenAI
response = await self.openai.chat.completions.create(
model="gpt-4o",
messages=payload["messages"],
max_tokens=2000
)
# Lagre resultat
await self._store_result(
payload["request_id"],
response.choices[0].message.content
)
# Callback hvis konfigurert
if payload.get("callback_url"):
await self._send_callback(
payload["callback_url"],
payload["request_id"],
response.choices[0].message.content
)
await receiver.complete_message(message)
except Exception as e:
if message.delivery_count < 3:
await receiver.abandon_message(message)
else:
await receiver.dead_letter_message(
message,
reason=str(e))
Azure Functions Queue Trigger
// Azure Function: Prosesser AI-forespørsler fra Storage Queue
using Azure.AI.OpenAI;
using Azure.Messaging.ServiceBus;
using Microsoft.Azure.Functions.Worker;
public class AIRequestProcessor
{
private readonly AzureOpenAIClient _openAIClient;
public AIRequestProcessor(AzureOpenAIClient openAIClient)
{
_openAIClient = openAIClient;
}
[Function("ProcessAIRequest")]
[ServiceBusOutput("ai-results", Connection = "ServiceBusConnection")]
public async Task<ServiceBusMessage> Run(
[ServiceBusTrigger("ai-requests",
Connection = "ServiceBusConnection")]
ServiceBusReceivedMessage message,
FunctionContext context)
{
var logger = context.GetLogger("ProcessAIRequest");
var request = JsonSerializer.Deserialize<AIRequest>(
message.Body.ToString());
logger.LogInformation(
"Processing request {RequestId}", request!.RequestId);
var chatClient = _openAIClient.GetChatClient("gpt-4o");
var response = await chatClient.CompleteChatAsync(
request.Messages.Select(m =>
new UserChatMessage(m.Content)).ToList());
var result = new AIResult
{
RequestId = request.RequestId,
Output = response.Value.Content[0].Text,
CompletedAt = DateTime.UtcNow,
TokensUsed = response.Value.Usage.TotalTokenCount
};
return new ServiceBusMessage(
JsonSerializer.Serialize(result))
{
MessageId = request.RequestId,
Subject = "completed"
};
}
}
Event-Driven Design
Azure OpenAI med Event Grid
# Event-driven pattern: Trigger AI-prosessering fra dokumenter
# Ny blob → Event Grid → Function → OpenAI → Result store
from azure.functions import Blueprint, EventGridEvent
from openai import AzureOpenAI
import json
bp = Blueprint()
@bp.event_grid_trigger(arg_name="event")
@bp.cosmos_db_output(
arg_name="resultDoc",
database_name="ai-results",
container_name="completions",
connection="CosmosConnection"
)
async def process_document_event(
event: EventGridEvent,
resultDoc: func.Out[str]
):
"""Process document when uploaded to Blob Storage."""
data = event.get_json()
blob_url = data["url"]
# Hent dokumentinnhold
document_text = await download_and_extract(blob_url)
# Prosesser med Azure OpenAI
client = AzureOpenAI(
azure_endpoint=os.environ["AZURE_OPENAI_ENDPOINT"],
api_key=os.environ["AZURE_OPENAI_KEY"],
api_version="2024-10-21"
)
response = client.chat.completions.create(
model="gpt-4o",
messages=[
{"role": "system", "content": "Analyser dette dokumentet..."},
{"role": "user", "content": document_text[:128000]}
],
max_tokens=2000
)
result = {
"id": event.id,
"source_blob": blob_url,
"analysis": response.choices[0].message.content,
"tokens_used": response.usage.total_tokens,
"processed_at": datetime.utcnow().isoformat()
}
resultDoc.set(json.dumps(result))
Request-Response Decoupling
Background Tasks med Azure OpenAI Responses API
from openai import AzureOpenAI
import time
def submit_background_task(client: AzureOpenAI, prompt: str) -> str:
"""Submit long-running task using background mode."""
response = client.responses.create(
model="o3", # Reasoning modell — kan ta minutter
input=prompt,
background=True # Kjør asynkront
)
return response.id
def poll_for_result(
client: AzureOpenAI,
response_id: str,
max_wait_seconds: int = 600,
poll_interval: int = 5
) -> dict:
"""Poll for background task completion."""
start = time.time()
while time.time() - start < max_wait_seconds:
result = client.responses.retrieve(response_id)
if result.status == "completed":
return {
"status": "completed",
"output": result.output,
"duration_seconds": round(time.time() - start, 1)
}
elif result.status == "failed":
return {"status": "failed", "error": result.error}
time.sleep(poll_interval)
return {"status": "timeout"}
# Bruk: Kompleks analyse som kan ta flere minutter
response_id = submit_background_task(
client,
"Analyser dette reguleringsverket og identifiser alle krav..."
)
# Klienten kan gjøre andre ting mens vi venter
result = poll_for_result(client, response_id)
Status Polling and Webhooks
Webhook-basert notifikasjon
# Webhook handler for Azure OpenAI events
from flask import Flask, request, Response
import hmac
import hashlib
app = Flask(__name__)
WEBHOOK_SECRET = os.environ["OPENAI_WEBHOOK_SECRET"]
@app.route("/webhooks/openai", methods=["POST"])
def handle_openai_webhook():
"""Handle Azure OpenAI webhook events."""
# Verifiser signatur
signature = request.headers.get("Webhook-Signature")
webhook_id = request.headers.get("Webhook-ID")
if not verify_signature(request.data, signature):
return Response("Invalid signature", status=400)
# Idempotency check
if is_already_processed(webhook_id):
return Response(status=200)
event = request.get_json()
# Prosesser event
if event.get("type") == "batch.completed":
handle_batch_complete(event["data"])
elif event.get("type") == "fine_tuning.job.succeeded":
handle_finetuning_complete(event["data"])
mark_as_processed(webhook_id)
return Response(status=200)
def verify_signature(payload: bytes, signature: str) -> bool:
"""Verify webhook signature."""
expected = hmac.new(
WEBHOOK_SECRET.encode(),
payload,
hashlib.sha256
).hexdigest()
return hmac.compare_digest(expected, signature)
# Polling-basert status-sjekk med exponential backoff
import asyncio
async def poll_with_backoff(
check_fn,
initial_interval: float = 2.0,
max_interval: float = 60.0,
backoff_factor: float = 1.5,
timeout: float = 3600.0
) -> dict:
"""Poll with exponential backoff until completion or timeout."""
interval = initial_interval
elapsed = 0.0
while elapsed < timeout:
result = await check_fn()
if result.get("status") in ("completed", "failed"):
return result
await asyncio.sleep(interval)
elapsed += interval
interval = min(interval * backoff_factor, max_interval)
return {"status": "timeout", "elapsed": elapsed}
REST API for Status Polling
// ASP.NET Core: Status polling endpoint for async AI requests
[ApiController]
[Route("api/ai")]
public class AIRequestController : ControllerBase
{
private readonly ICosmosDbService _cosmosDb;
private readonly IServiceBusSender _sender;
[HttpPost("requests")]
public async Task<IActionResult> SubmitRequest(
[FromBody] AIRequestDto request)
{
var requestId = Guid.NewGuid().ToString();
// Legg i kø for asynkron prosessering
await _sender.SendAsync(new ServiceBusMessage(
JsonSerializer.Serialize(request))
{
MessageId = requestId
});
// Returner 202 Accepted med Location header
return AcceptedAtAction(
nameof(GetStatus),
new { requestId },
new { requestId, status = "queued" });
}
[HttpGet("requests/{requestId}/status")]
public async Task<IActionResult> GetStatus(string requestId)
{
var result = await _cosmosDb.GetRequestStatus(requestId);
if (result == null)
return NotFound();
if (result.Status == "completed")
return Ok(result);
// Returnér 200 med status og Retry-After header
Response.Headers.Append("Retry-After", "5");
return Ok(new { requestId, status = result.Status });
}
}
Event-Driven Architecture Styles (oppdatert 2026-04)
Microsoft dokumenterer to primære topologier for event-drevet AI-prosessering:
Broker-topologi vs. Mediator-topologi
| Aspekt | Broker-topologi | Mediator-topologi |
|---|---|---|
| Koordinering | Events publiseres direkte til broker | Central mediator koordinerer workflow |
| Eksempel | Azure Event Hubs + Service Bus | Azure Durable Functions |
| Kobling | Løs kobling mellom produsenter/konsumenter | Sterkere kobling via mediator |
| Bruksscenario | Høyvolum streaming, uavhengige konsumenter | Komplekse AI-arbeidsflyter med avhengigheter |
Azure Event Hubs vs. Azure Event Grid
| Service | Type | Bruksscenario |
|---|---|---|
| Azure Event Hubs | Durable event stream (log) | AI-inferensresultater som skal prosesseres av mange konsumenter |
| Azure Event Grid | Publish-subscribe, reaktiv | Trigger AI-jobb ved filnedlasting, blob-endring |
| Azure Service Bus | Message queue, garantert levering | Jobb-kø for AI-prosessering med retry og dead-letter |
Utfordringer i event-drevne AI-arkitekturer
# Utfordring 1: Garantert levering
# Bruk Service Bus med peek-lock for å garantere at AI-jobb fullføres
from azure.servicebus import ServiceBusClient, ServiceBusMessage
import json
def process_ai_job_safely(
servicebus_conn: str,
queue_name: str,
ai_processor
) -> None:
"""Garantert levering via peek-lock mønster."""
with ServiceBusClient.from_connection_string(servicebus_conn) as sb:
with sb.get_queue_receiver(queue_name, max_wait_time=5) as receiver:
for message in receiver:
# Peek-lock: meldingen er reservert, ikke slettet
try:
payload = json.loads(str(message))
result = ai_processor(payload)
# Fullfør melding (slett fra kø) kun ved suksess
receiver.complete_message(message)
publish_result(result)
except Exception as e:
# Abandon: meldingen returneres til kø for ny levering
receiver.abandon_message(message)
# Utfordring 2: Eventual consistency
# AI-resultater publiseres asynkront — bruk correlation ID for sporing
def create_ai_job(correlation_id: str, payload: dict) -> dict:
"""Returner job receipt umiddelbart, resultat kommer asynkront."""
return {
"correlation_id": correlation_id,
"status": "accepted",
"result_url": f"/api/results/{correlation_id}",
"estimated_completion_seconds": 30
}
# Utfordring 3: Ordregaranti
# Event Hubs garanterer ordre innen én partisjon
# Bruk samme partisjonsnøkkel for relaterte AI-forespørsler
def publish_ordered_event(
producer,
partition_key: str, # f.eks. dokument-ID
event_data: dict
) -> None:
from azure.eventhub import EventData
event = EventData(json.dumps(event_data))
event.properties = {"partition_key": partition_key}
producer.send_batch([event], partition_key=partition_key)
Norsk offentlig sektor
- Saksbehandlingssystemer: Asynkron prosessering er ideelt for AI-assistert saksbehandling der analyse kan ta tid. Saksbehandler sender inn dokument, fortsetter med annet arbeid, og mottar notifikasjon når analysen er ferdig.
- Arkivloven: Sørg for at alle mellomliggende meldinger i køer (Service Bus, Queue Storage) krypteres og at sensitive data ikke lagres utover nødvendig prosesseringstid.
- Personvern: Dead letter queues kan inneholde personopplysninger — konfigurer automatisk sletting og monitorering av DLQ-dybde.
- Tilgjengelighet: Asynkrone mønstre forbedrer brukeropplevelsen for tjenester med krav om universell utforming — brukere slipper å vente på skjermen.
- Batch-prosessering: Bruk Azure OpenAI Batch API for periodiske oppgaver (nattlige rapporter, ukentlige analyser) med 50% kostnadsreduksjon.
Beslutningsrammeverk
| Scenario | Anbefaling | Begrunnelse |
|---|---|---|
| Bruker venter på svar (<3s) | Synkron + streaming | Best brukeropplevelse for korte svar |
| Dokumentanalyse (minutter) | Service Bus kø + polling | Bruker kan gjøre annet arbeid |
| Reasoning-modell (o3/o1) | Background Tasks API | Innebygd asynkron prosessering |
| Stort batch-volum (1000+) | Azure OpenAI Batch API | 50% kostnadsreduksjon |
| Event-drevet pipeline | Event Grid + Functions | Automatisk trigger ved nye data |
| Kritisk pålitelighet | Service Bus + DLQ | Garantert leveranse og feilhåndtering |
Referanser
- Azure OpenAI Batch API — Batch processing
- Azure OpenAI Responses API — Background tasks — Background mode
- Azure OpenAI Webhooks — Event notifications
- Event-driven architecture style — Architecture patterns
- Azure Functions on Container Apps — Event-driven compute
For Cosmo
- Bruk denne referansen når kunden har AI-workloads som ikke krever umiddelbart svar, eller når de opplever timeout-problemer med langvarige AI-forespørsler.
- Azure OpenAI Background Tasks er den enkleste løsningen for reasoning-modeller (o3, o1) som kan ta minutter — sett
background: true. - For enterprise-arkitekturer, anbefal Service Bus fremfor Queue Storage — gir sessions, dead letter queues og transaksjonsstøtte.
- Implementer alltid idempotency i webhook-handlere og consumers — meldinger kan leveres mer enn én gang.
- Batch API bør være standard for alle ikke-sanntids workloads — 50% kostnadsreduksjon er en enkel gevinst.